بيئة تشغيل الوكلاء المُدارة

الاستقرار: تجريبي — هذه الميزة إضافية ومقيّدة بميزة خلف managed-runtime. ولا تؤثر في Runner/LlmAgent الحالية أو APIs عند تعطيل الميزة. وقد تتغير واجهة API في الإصدارات المستقبلية.

نظرة عامة

إن بيئة تشغيل الوكلاء المُدارة (adk-managed) هي محرك لتنفيذ الوكلاء، متين وقابل للاستئناف ومحايد تجاه المزوّد. إذ يتلقى ManagedAgentDef تصريحيًا، وينشئ وكيلًا قابلًا للتشغيل، ويديره باعتباره جلسة خلفية قابلة للاستئناف من نقاط التحقق وتدفق الأحداث.

إن بيئة التشغيل هي مكتبة وليست خدمة. وتستضيفها المنصة. وهذا يعني:

  • قابلة للاختبار بمعزل عن غيرها: لا توجد تبعيات على HTTP/المصادقة/الفوترة
  • قابلة للتضمين: تستخدم عمليات النشر المستضافة ذاتيًا سمة بيئة التشغيل نفسها مباشرةً
  • منصة قابلة للاستبدال: يمكن لمنصات مختلفة استضافة بيئة التشغيل نفسها
  • محايدة تجاه المزوّد: تسلسلات أحداث متطابقة بغض النظر عن مزوّد النموذج

البدء السريع

أضف الميزة إلى Cargo.toml:

[dependencies]
adk-rust = { version = "2.1.0", features = ["managed-runtime"] }

أو استخدم الحزمة adk-managed مباشرةً:

[dependencies]
adk-managed = "2.1.0"
adk-session = "2.1.0"

مثال أدنى (ScriptedLlm — بدون مفتاح API)

use std::sync::Arc;
use adk_managed::{
    DefaultManagedAgentRuntime, ManagedAgentRuntime, ModelResolver,
    ScriptedLlm, ScriptedTurn,
    resolver::ResolverResult,
    types::{ContentBlock, ManagedAgentDef, ModelRef, UserEvent},
};
use adk_session::InMemorySessionService;
use async_trait::async_trait;
use futures::StreamExt;

// A resolver that returns our scripted LLM
struct MockResolver { llm: Arc<dyn adk_core::Llm> }

#[async_trait]
impl ModelResolver for MockResolver {
    async fn resolve(&self, _: &ModelRef) -> ResolverResult<Arc<dyn adk_core::Llm>> {
        Ok(self.llm.clone())
    }
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 1. Create a scripted LLM (deterministic, offline, $0)
    let llm = Arc::new(ScriptedLlm::new("test-model", vec![
        ScriptedTurn { text: Some("Hello!".into()), tool_calls: vec![] },
    ]));

    // 2. Build the runtime
    let runtime = DefaultManagedAgentRuntime::new(
        Arc::new(MockResolver { llm }),
        Arc::new(InMemorySessionService::new()),
    );

    // 3. Create an agent
    let def = ManagedAgentDef::new("my-agent", ModelRef::Shorthand("test-model".into()))
        .with_system("You are helpful.");
    let agent = runtime.create(def).await?;

    // 4. Start a session
    let session = runtime.start_session(&agent, None).await?;

    // 5. Subscribe to events and send a message
    let mut stream = runtime.stream_events(&session, None).await?;
    runtime.send_event(&session, UserEvent::Message {
        content: vec![ContentBlock::Text { text: "Hi!".into() }],
    }).await?;

    // 6. Collect events
    while let Some(event) = stream.next().await {
        println!("{event:?}");
    }
    Ok(())
}

البنية

┌─────────────────────────────────────────────────────────────┐
│                Platform Layer (ep-* crates)                  │
│    HTTP Routes │ Auth │ Billing │ Multi-tenancy              │
└──────────────────────────┬──────────────────────────────────┘
                           │ Rust trait calls (in-process)
                           ▼
┌─────────────────────────────────────────────────────────────┐
│            Runtime Layer (adk-managed)                       │
│                                                             │
│  ManagedAgentRuntime trait + DefaultManagedAgentRuntime      │
│  ───────────────────────────────────────────────────        │
│  • Builds runnable agents from ManagedAgentDef              │
│  • Runs supervised session loop (durable, resumable)        │
│  • Emits provider-neutral SessionEvent stream               │
│  • Manages custom tool parking, checkpoints, interrupts     │
│  • Resolves ModelRef → Arc<dyn Llm>                         │
│                                                             │
│  Composes existing crates:                                  │
│  ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐    │
│  │adk-runner│ │adk-session│ │adk-model │ │adk-tool    │    │
│  └──────────┘ └──────────┘ └──────────┘ └────────────┘    │
└─────────────────────────────────────────────────────────────┘

الأنواع الأساسية

سمة ManagedAgentRuntime

السمة غير المتزامنة المركزية التي تحدد دورة حياة الوكيل كاملةً:

الطريقةالوصف
create(def)تسجيل تعريف وكيل، يعيد AgentHandle
start_session(agent, env?)بدء جلسة جديدة، الحالة الأولية Queued
send_event(session, event)أرسل UserEvent إلى الجلسة
stream_events(session, from_seq?)اشترك في تدفق SessionEvent
interrupt(session)توقّف عند الحد التالي، وأصدر status.idle
pause(session)أنشئ نقطة تحقق وأوقف المعالجة
resume(session)الاستئناف من الإيقاف المؤقت أو إعادة التشغيل
status(session)الاستعلام عن SessionStatus الحالي
archive(session)حالة نهائية، مع الاحتفاظ بالبيانات
delete_session(session)إزالة بيانات الجلسة

ManagedAgentDef

تعريف وكيل تصريحي باستخدام أداة البناء API:

let def = ManagedAgentDef::new("my-agent", ModelRef::Shorthand("gemini-3.7-flash".into()))
    .with_system("You are a helpful assistant.")
    .with_description("Research agent with web search")
    .with_tools(vec![ToolConfig::BuiltIn(ManagedBuiltinTool::WebSearch)]);

SessionEvent

تدفق أحداث محايد تجاه موفّر الخدمة مع أرقام تسلسلية رتيبة:

  • agent.message — محتوى نص المساعد
  • agent.tool_use — استدعاء أداة مضمّنة
  • agent.custom_tool_use — أداة مخصّصة ينفّذها العميل (توقّف الحلقة)
  • agent.mcp_tool_use — استدعاء أداة MCP
  • status.running — بدء الدورة
  • status.idle — اكتمال الدورة (مع stop_reason)
  • error — خطأ في التنفيذ

UserEvent

أحداث من العميل إلى الوكيل:

  • user.message — إرسال محتوى إلى الوكيل
  • user.interrupt — إيقاف الدورة الحالية
  • user.tool_confirmation — السماح بتنفيذ الأداة أو رفضه
  • user.custom_tool_result — إرجاع نتائج الأداة المخصّصة
  • user.tool_result — نتيجة الأداة المضمّنة (مستضاف ذاتيًا فقط)
  • user.define_outcome — تعيين معايير النجاح

ModelRef

مرجع نموذج محايد تجاه موفّر الخدمة ويدعم جميع الموفّرين:

// Shorthand (provider inferred from name)
ModelRef::Shorthand("gemini-3.7-flash".into())
ModelRef::Shorthand("gpt-5.6-terra".into())
ModelRef::Shorthand("claude-sonnet-5".into())

// Structured (explicit provider)
ModelRef::Structured {
    provider: Provider::OpenaiCompatible,
    model: ModelConfig::Compatible {
        model: "my-model".into(),
        base_url: "https://my-endpoint.com/v1".into(),
        api_key: "sk-...".into(),
    },
    speed: None,
}

الميزات الرئيسية

الجلسات المستديمة

يتم إنشاء نقطة تحقق ذرّية لكل حدث. عند تعطل العملية، يعيد resume() تكوين حالته من آخر نقطة تحقق متسقة دون فقدان أي أحداث:

// Before crash: events 0..5 committed
// After restart:
runtime.resume(&session).await?;
// Continues from seq=5, no gap, no duplicate

إيقاف الأدوات المخصّصة مؤقتًا

عندما يُصدر الوكيل agent.custom_tool_use، تتوقف الحلقة مؤقتًا إلى أن يعيد العميل النتائج أو تنقضي مهلة قابلة للتهيئة:

// Agent emits: agent.custom_tool_use { custom_tool_use_id: "ct_1", name: "deploy" }
// Client executes the tool, then:
runtime.send_event(&session, UserEvent::CustomToolResult {
    custom_tool_use_id: "ct_1".into(),
    content: vec![ContentBlock::Text { text: "Deployed successfully".into() }],
}).await?;

إعادة تشغيل الأحداث

دعم إعادة اتصال SSE Last-Event-ID عبر إعادة التشغيل المستندة إلى التسلسل:

// Reconnect from seq 42 — replays events 43, 44, ... then live tail
let stream = runtime.stream_events(&session, Some(42)).await?;

تكافؤ الموفّرين

ينتج ManagedAgentDef المتطابق تسلسلات متطابقة بايتًا لأنواع الأحداث عبر Gemini وOpenAI وAnthropic وOllama والموفّرين المتوافقين مع OpenAI (التركيبة F-8).

الاختبار باستخدام ScriptedLlm

إن ScriptedLlm هو بديل LLM حتمي يختبر مسار التنفيذ الكامل. يتم استبدال استدعاء الموفّر API فقط:

use adk_managed::testing::{ScriptedLlm, ScriptedTurn, ScriptedToolCall};
use serde_json::json;

let llm = ScriptedLlm::new("test", vec![
    ScriptedTurn {
        text: Some("I'll search for that.".into()),
        tool_calls: vec![ScriptedToolCall {
            name: "web_search".into(),
            input: json!({"query": "rust agents"}),
            id: Some("tc_1".into()),
        }],
    },
    ScriptedTurn {
        text: Some("Here are the results...".into()),
        tool_calls: vec![],
    },
]);

مرجع API

تتوفر وثائق API الكاملة على docs.rs:

مثال على اختبار الدخان

يتوفر crate مستقل نموذجي لفرق المنصة:

cargo run --manifest-path examples/managed_runtime_hello/Cargo.toml

يُشغّل هذا الاختبار التجهيزي F-1 من البداية إلى النهاية باستخدام ScriptedLlm (ولا يتطلب مفتاح API).

متانة حالة الجلسة المُدارة

تُخزَّن حالة الجلسة المُدارة — سجل الأحداث، وموضع التسلسل، واستدعاءات الأدوات المتوقفة مؤقتًا، وحالة دورة الحياة — في ManagedStateStore. ويُبلغ المخزن عن ضمانه الخاص:

المتانةالمعنى
ProcessLocalتعمل إعادة التشغيل والاستئناف أثناء تشغيل العملية. يؤدي التعطل إلى فقدان الحالة، ولا يمكن لعملية أخرى استئناف الجلسة.
CrashDurableتُكتب الحالة إلى مخزن داعم قبل الإقرار بالكتابة، بحيث يمكن لعملية أخرى إعادة إنشاء الجلسة.

لا يتم شحن سوى InMemoryManagedStateStore، وهو ProcessLocal. تحقّق من الضمان بدلًا من استنتاجه من وجود حفظ نقاط التحقق:

use adk_managed::{Durability, InMemoryManagedStateStore, ManagedStateStore};

let store = InMemoryManagedStateStore::new();
assert_eq!(store.durability(), Durability::ProcessLocal);
assert!(!store.durability().survives_process_loss());

حفظ نقاط التحقق مقابل التفريغ

يسجّل CheckpointManager::checkpoint حدثًا وحالة التشغيل الجديدة معًا، لذا لا يرى الاست replay أحدهما دون الآخر. هذه عملية كتابة في حقول المدير نفسه. يكتب flush اللقطة في المخزن المُكوَّن، ويعيد restore بناء مدير منها:

use adk_managed::{CheckpointManager, InMemoryManagedStateStore, ManagedStateStore};
use std::sync::Arc;

# async fn example() -> Result<(), adk_managed::types::RuntimeError> {
let store: Arc<dyn ManagedStateStore> = Arc::new(InMemoryManagedStateStore::new());
let manager = CheckpointManager::new("session-1".to_string()).with_store(Arc::clone(&store));
manager.flush().await?;

let restored = CheckpointManager::restore("session-1".to_string(), store).await?;
assert_eq!(restored.session_id(), "session-1");
# Ok(())
# }

الإبلاغ عن الحالة

يقرأ ManagedAgentRuntime::status المقبض نفسه الذي تكتب إليه حلقة الجلسة، لذا تكون الانتقالات العادية مرئية، وليس فقط انتقالات مستوى التحكم:

الانتقالالسبب
QueuedRunningتبدأ جولة
RunningIdleتكتمل الجولة ويتم تسجيل الاستخدام
أيّ → Pausedpause
أيّ → Archivedarchive أو delete_session

ملاحظة: قبل ذلك، كان هذا مقبضًا مشتركًا واحدًا، وكان status يُبلّغ عن Queued طوال مدة الجلسة، بما في ذلك أثناء تنفيذها للدورات. وكانت انتقالات مستوى التحكم (الإيقاف المؤقت، والاستئناف، والأرشفة) مرئية لأنها كانت تكتب مباشرةً إلى المقبض.

دلالات الحذف

يُزيل delete_session كلا المستويين:

  1. يضع الجلسة في الحالة النهائية ويلغي حلقتها.
  2. يزيل مقبض وقت التشغيل.
  3. يحذف المحادثة المستمرة من خلال SessionService المُحقن، باستخدام الهوية نفسها start_session التي أنشأتها.

إذا فشلت الخطوة 3، يُرجع delete_session خطأً يذكر التطبيق والمستخدم والجلسة التي لا تزال تحتفظ بالبيانات — إذ يكون المقبض قد أُزيل بالفعل في تلك المرحلة، ولذلك يجب إبلاغ المستدعي بما يحتاج إلى تنظيف يدوي بدلًا من السماح له بافتراض نجاح العملية.

runtime.delete_session(&session).await?;
// The handle is gone and the conversation is no longer in the session backend.

مهم: يحذف الحذفُ المحادثة الخاصة بالجلسة لدى مالكها. راجع ملكية الجلسة.

ملكية الجلسة

يتطلب start_session وجود ManagedOwner. تُحفظ الجلسة تحت تلك الهوية، وتستخدمها كل استدعاءات Runner التي تجريها حلقة الجلسة:

use adk_managed::{ManagedAgentRuntime, ManagedOwner};

# async fn start(runtime: &dyn ManagedAgentRuntime, agent: &adk_managed::AgentHandle)
# -> Result<(), adk_managed::RuntimeError> {
let owner = ManagedOwner::new("support-console", "user-42")?;
let session = runtime.start_session(agent, &owner, None).await?;
# Ok(())
# }

مهم: وُثّق checkpoint على أنه «يحفظ بشكل ذري» مع ضمان أن «إعادة التشغيل ستعرض رؤية متسقة بعد أي تعطل»، كما وُصف التحميل بأنه يعيد «كل ما يلزم لإعادة بناء جلسة بعد إعادة التشغيل». ولم يكن أيٌّ من ذلك صحيحًا: فكلاهما كان يعمل على حقول موجودة في الذاكرة فقط، من دون معاملة على أي مخزن مستمر. ومع المخزن المشحون، لا يجد restore شيئًا في عملية جديدة. كلا المكوّنين مطلوب، ويجب ألا يكون أيٌّ منهما فارغًا. وتُعنون الجلسات التابعة لمالكين مختلفين بشكل منفصل، لذا يقتصر نطاق البحث والحذف على مالك واحد ولا يمكنهما الوصول إلى بيانات مالك آخر.

ملاحظة: كانت كل جلسة مُدارة تُحفَظ سابقًا ضمن الثابتين managed / managed_user، ولذلك كانت جميعها تشترك في مساحة أسماء منطقية واحدة: لم يكن من الممكن تحديد نطاق أي شيء ليقتصر على مستدعٍ، ولم يكن بالإمكان إسناد أي جلسة إلى مستدعٍ بعينه.

إعداد البيئة

يحمل EnvironmentConfig كلًا من env_vars وworking_dir. يرفض وقت التشغيل هذا أي إعداد يطلب أيًا منهما:

invalid request: EnvironmentConfig cannot be honoured by this runtime: sessions run
in-process, so per-session environment variables and working directories would have to mutate
process-global state shared with other sessions. Pass `None`, or configure a sandboxed runtime.

تُشغَّل الجلسات داخل العملية، ولذلك فإن تطبيق متغيرات بيئة خاصة بكل جلسة أو دليل عمل سيؤدي إلى تغيير الحالة المشتركة مع كل جلسة أخرى. إن الرفض هو النتيجة الصادقة؛ فحدّ تنفيذ معزول هو ما سيجعل الطلب قابلًا للتنفيذ.

ملاحظة: كان اسم الوسيط سابقًا _env ثم جرى تجاهله، ولذلك كان المستدعي الذي يزوّد إعدادات البيئة يحصل على جلسة تتجاهلها بصمت.

بيئة تشغيل الوكلاء المُدارة - وثائق ADK-Rust | ADK-Rust