fix agent idle stream heartbeat visibility

This commit is contained in:
Hamza Ayed
2026-10-04 16:37:45 +03:00
parent cfd15e9928
commit 23f9ccc99d
4 changed files with 50 additions and 3 deletions
+2
View File
@@ -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/الرابط الرمزي؛ لم يتغير كاش الحزمة أو مصدرها.
+3 -3
View File
@@ -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"
@@ -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;
@@ -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()