From 23f9ccc99dd452ab33c36467f01a4859237f46f6 Mon Sep 17 00:00:00 2001 From: Hamza Ayed Date: Sun, 4 Oct 2026 16:37:45 +0300 Subject: [PATCH] fix agent idle stream heartbeat visibility --- SovereignAI-Starter/ROADMAP.md | 2 + SovereignAI-Starter/app/main.py | 6 +-- .../lib/core/network/api_repository.dart | 3 ++ .../tests/test_agent_stream.py | 42 +++++++++++++++++++ 4 files changed, 50 insertions(+), 3 deletions(-) create mode 100644 SovereignAI-Starter/tests/test_agent_stream.py diff --git a/SovereignAI-Starter/ROADMAP.md b/SovereignAI-Starter/ROADMAP.md index 487a9c4..79b7909 100644 --- a/SovereignAI-Starter/ROADMAP.md +++ b/SovereignAI-Starter/ROADMAP.md @@ -6,6 +6,8 @@ - 2026-10-04 — إصلاح سبب انتهاء مهلة الوكيل بعد عدة دقائق: بعد إضافة SSE keep-alive، كشف الفحص أن كل استدعاء نموذج داخل حلقة الوكيل ما زال يستخدم المهلة الافتراضية 180 ثانية، رغم أن العميل والبث يسمحان حتى 10 دقائق. رُفعت مهلة كل خطوة نموذج في الوكيل إلى 600 ثانية لتتسق مع مهلة البث، وأضيف اختبار يثبت قيمة المهلة. بعد إعادة تشغيل Uvicorn للمشروع على 8000، نجحت اختبارات `test_timeouts_and_cancellation` 5/5، وأعاد `/health` الحالة `ok`. طلب وكيل حي أعاد HTTP 200 واكتمل بعد 41.9 ثانية مع نبضتي keep-alive؛ طلب حي ثانٍ اكتمل بعد 12.2 ثانية مع أحداث `done` وإغلاق الجلسة. لا يثبت هذا اختبارًا يستغرق خمس دقائق فعليًا، لكنه يزيل حد الـ180 ثانية الذي كان يطابق الانقطاع المبلغ عنه؛ ما زال اختبار نافذة Flutter التفاعلي غير منجز. +- 2026-10-04 — متابعة بلاغ انقطاع الوكيل بعد 4–5 دقائق: `/health` على `127.0.0.1:8000` سليم ويستخدم `gemma4:e2b`. نفذت طلبًا حيًا عبر `/v1/agent/run/stream` مع `app/main.py`؛ استغرق 90.2 ثانية، وأعاد HTTP 200 ثم أحداث التقدم وkeep-alive وانتهى بـ`done` ومعرّفي تتبع للطلب والتدقيق. لم يتكرر الانقطاع في هذه التجربة الأقصر من البلاغ. حُوّلت نبضة SSE من تعليق صامت إلى `event: heartbeat` مرئي يعرض للواجهة أن الوكيل ما زال يعمل؛ اختبار انتظار خامل اصطناعي نجح (1/1) وأكد heartbeat ثم النتيجة النهائية. يلزم تحقق Flutter وتكرار اختبار يستمر 4–5 دقائق لتحديد سبب البلاغ الأصلي؛ هذه النتيجة لا تثبت بعد زواله. + - 2026-10-04 — إصلاح بقاء اتصال بث الوكيل: ظهر أن SSE لا يرسل أي بايت أثناء انتظار استجابة نموذج بطيئة، مع أن مهلة Flutter عشر دقائق؛ أضيف تعليق `: keep-alive` كل 15 ثانية عند خمول البث. اختبار حتمي يحاكي انتظار النموذج وإلغاء العميل، ومجموعة `test_timeouts_and_cancellation.py` نجحت 4/4. أُعيد تشغيل API على 8000 من الشفرة الحالية؛ `/health` أعاد 200، وطلب وكيل حي اختار الحاسبة وأعاد `5.0`، وظهر تعليقان keep-alive قبل اكتمال الرد ثم نجح تسجيل الخروج. هذا يعالج إغلاق الاتصالات الخاملة على المسارات الوسيطة، ولا يثبت سرعة النموذج نفسه. - 2026-10-04 — إصلاح بطء فهرسة PDF الممسوح الذي كان يوقف حلقة FastAPI: نقل تحليل الصفحات والرسم وOCR إلى worker threads؛ صار فهرس Flutter يرسل ملفًا واحدًا في كل مرة بمهلة 5 دقائق ويعرض اسم الملف ورقمه. أثناء OCR حي استغرق 155.21 ثانية، وظل `/health` يستجيب خلال 0.02 ثانية، واكتمل OCR للصفحة وفهرستها دلاليًا. اجتاز اختبار PDF المختلط 1/1؛ واجتازت Flutter 22/22 و`flutter analyze --no-pub` بلا ملاحظات، واختبارات مهلة الوكيل 4/4. لم يكتمل بناء Windows Debug بهذه التغييرات في هذه الجولة: نسخة التطوير المؤقتة تعطلت بسبب مسار تضمين `flutter_secure_storage_windows` في MSBuild/الرابط الرمزي؛ لم يتغير كاش الحزمة أو مصدرها. diff --git a/SovereignAI-Starter/app/main.py b/SovereignAI-Starter/app/main.py index 71b1425..0374d79 100644 --- a/SovereignAI-Starter/app/main.py +++ b/SovereignAI-Starter/app/main.py @@ -2385,9 +2385,9 @@ async def run_agent_stream( message = await asyncio.wait_for(queue.get(), timeout=0.25) except TimeoutError: if perf_counter() - last_output_at >= AGENT_STREAM_HEARTBEAT_SECONDS: - # Keep otherwise-idle SSE connections alive while the model - # or a local tool is taking several minutes to finish. - yield ": keep-alive\n\n" + # Send a visible SSE event (not only a comment) so clients + # can reset idle timers and tell the user the agent is alive. + yield "event: heartbeat\ndata: {}\n\n" last_output_at = perf_counter() continue yield f"event: progress\ndata: {json.dumps({'message': message}, ensure_ascii=False)}\n\n" diff --git a/SovereignAI-Starter/flutter_app/lib/core/network/api_repository.dart b/SovereignAI-Starter/flutter_app/lib/core/network/api_repository.dart index 0caae1e..72df7bf 100644 --- a/SovereignAI-Starter/flutter_app/lib/core/network/api_repository.dart +++ b/SovereignAI-Starter/flutter_app/lib/core/network/api_repository.dart @@ -587,6 +587,9 @@ class ApiRepository { final message = data['message']?.toString(); if (message != null) onProgress?.call(message); break; + case 'heartbeat': + onProgress?.call('الوكيل ما زال يعمل، أنتظر النموذج المحلي...'); + break; case 'done': resultData = data; break; diff --git a/SovereignAI-Starter/tests/test_agent_stream.py b/SovereignAI-Starter/tests/test_agent_stream.py new file mode 100644 index 0000000..690fdf3 --- /dev/null +++ b/SovereignAI-Starter/tests/test_agent_stream.py @@ -0,0 +1,42 @@ +import asyncio +import os +import tempfile +import unittest +from unittest.mock import patch + +_TEST_DATA_DIR = None +if "SOVEREIGNAI_DATA_DIR" not in os.environ: + _TEST_DATA_DIR = tempfile.TemporaryDirectory(prefix="sovereignai-stream-tests-") + os.environ["SOVEREIGNAI_DATA_DIR"] = _TEST_DATA_DIR.name + +from app import main + + +class AgentStreamTests(unittest.IsolatedAsyncioTestCase): + async def test_idle_model_wait_sends_heartbeat_then_final_result(self) -> None: + async def slow_agent(request, report_progress=None, *, user_id): + await report_progress("بدأ تحليل المهمة.") + await asyncio.sleep(0.05) + return {"task": request.task, "result": "اكتمل التحليل."} + + with ( + patch.object(main, "AGENT_STREAM_HEARTBEAT_SECONDS", 0.01), + patch.object(main, "_execute_agent", side_effect=slow_agent), + ): + response = await main.run_agent_stream( + main.AgentRequest(task="اختبار انتظار النموذج"), user_id="test-user" + ) + chunks = [chunk async for chunk in response.body_iterator] + + body = b"".join( + chunk if isinstance(chunk, bytes) else chunk.encode("utf-8") + for chunk in chunks + ).decode("utf-8") + self.assertIn("event: progress", body) + self.assertIn("event: heartbeat\ndata: {}", body) + self.assertIn("event: done", body) + self.assertIn("اكتمل التحليل", body) + + +if __name__ == "__main__": + unittest.main()