-
Notifications
You must be signed in to change notification settings - Fork 8.2k
Lorenze/feat/conversational flows #5896
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from 11 commits
Commits
Show all changes
30 commits
Select commit
Hold shift + click to select a range
9e8b47f
feat: add conversational flows documentation and chat session support
lorenzejay eca18b0
linted
lorenzejay f1c5ea3
feat: enhance flow event tracing and session management
lorenzejay 9f0292c
updated docs
lorenzejay 22f8fa6
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay f2254d9
feat: introduce experimental conversational flow framework
lorenzejay d4882c6
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay e71b916
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay 1bf73b2
handled docs
lorenzejay 38e762d
feat(flow): enhance conversational flow handling and tracing
lorenzejay b8237f6
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay ff46625
fix multimodal test
lorenzejay d8c8ab9
better conversational
lorenzejay 4d2909b
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay 1dd3af0
adjusted prompt
lorenzejay 796ac69
drop unused
lorenzejay 5f180de
fix test
lorenzejay b0a6ac1
refactor: rename to and update related documentation
lorenzejay 13f8713
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay 9b3900b
fix test
lorenzejay 14caf0b
merge conflict resolution
lorenzejay 9d286c4
adding experimetnal indicators
lorenzejay 5854e2b
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay f93014f
fix test and reloaded cassettes
lorenzejay 7fa1522
cleanup ConversationalFlow class
lorenzejay cb31d43
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay d6f6973
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay 340ec75
addressing double finalization and fixed tests
lorenzejay d474827
improve on emphemeral tracing and adddressing comments
lorenzejay 5111d80
Merge branch 'main' of github.com:crewAIInc/crewAI into lorenze/feat/…
lorenzejay File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,365 @@ | ||
| --- | ||
| title: تدفقات المحادثة | ||
| description: أنشئ تطبيقات دردشة متعددة الجولات مع kickoff لكل جولة وسجل الرسائل وتوجيه النية والتتبع وجسور WebSocket. | ||
| icon: comments | ||
| mode: "wide" | ||
| --- | ||
|
|
||
| ## نظرة عامة | ||
|
|
||
| تعامل التطبيقات المحادثية مع كل سطر من المستخدم كـ **تشغيل flow جديد** بنفس **معرّف الجلسة**. توفر CrewAI مساعدات لسجل الرسائل وتصنيف النية الاختياري وتأجيل التتبع وجسور الواجهة — دون API منفصل `chat()` على `Flow`. | ||
|
|
||
| | المفهوم | التنفيذ | | ||
| |---------|---------| | ||
| | معرّف الجلسة | `kickoff(session_id=...)` → `inputs["id"]` → `state.id` | | ||
| | سطر المستخدم | `kickoff(user_message=...)` يُضاف إلى `state.messages` قبل تشغيل الرسم | | ||
| | اكتمال الجولة | `FlowFinished` لهذا **التشغيل** فقط؛ تستمر المحادثة في `kickoff` التالي | | ||
| | تتبع الجلسة | `ConversationalConfig(defer_trace_finalization=True)` + `finalize_session_traces()` | | ||
|
|
||
| ## نقطة دخول واحدة: `kickoff` | ||
|
|
||
| استخدم **`flow.kickoff(user_message=..., session_id=...)`** لكل رسالة مستخدم (REST أو WebSocket أو CLI). لا تنشئ غلاف `chat()` مخصصاً على `Flow`. | ||
|
|
||
| | API | الاستخدام | | ||
| |-----|-----------| | ||
| | `kickoff(user_message=..., session_id=...)` | كل رسالة مستخدم | | ||
| | `kickoff_async(...)` | نفس المعاملات؛ دخول async أصلي | | ||
| | `ask()` | مطالبة حاجزة **داخل** خطوة واحدة | | ||
| | `@human_feedback` | الموافقة/الرفض على **مخرجات خطوة** — وليس السطر التالي | | ||
| | `ChatSession.handle_turn(...)` | طبقة نقل فوق `kickoff` | | ||
|
|
||
| ## بداية سريعة | ||
|
|
||
| ```python | ||
| from uuid import uuid4 | ||
|
|
||
| from crewai.flow import ( | ||
| ChatState, | ||
| ConversationalConfig, | ||
| Flow, | ||
| listen, | ||
| or_, | ||
| persist, | ||
| router, | ||
| start, | ||
| ) | ||
| from crewai.flow.persistence import SQLiteFlowPersistence | ||
|
|
||
|
|
||
| class SupportFlow(Flow[ChatState]): | ||
| conversational_config = ConversationalConfig( | ||
| default_intents=["order", "help", "goodbye"], | ||
| intent_llm="gpt-4o-mini", | ||
| defer_trace_finalization=True, | ||
| ) | ||
|
|
||
| @start() | ||
| def bootstrap(self): | ||
| if not self.state.session_ready: | ||
| self.state.session_ready = True | ||
| return "ready" | ||
|
|
||
| @router(bootstrap) | ||
| def route(self): | ||
| return self.state.last_intent or "help" | ||
|
|
||
| @listen("order") | ||
| def handle_order(self): | ||
| reply = "طلبك في الطريق." | ||
| self.append_message("assistant", reply) | ||
| return reply | ||
|
|
||
| @listen("help") | ||
| def handle_help(self): | ||
| reply = "كيف يمكنني المساعدة؟" | ||
| self.append_message("assistant", reply) | ||
| return reply | ||
|
|
||
| @listen("goodbye") | ||
| def handle_goodbye(self): | ||
| reply = "وداعاً!" | ||
| self.append_message("assistant", reply) | ||
| return reply | ||
|
|
||
| @persist(SQLiteFlowPersistence("support.db")) | ||
| @listen(or_(handle_order, handle_help, handle_goodbye)) | ||
| def finalize(self): | ||
| return self.state.model_dump() | ||
|
|
||
|
|
||
| session_id = str(uuid4()) | ||
| flow = SupportFlow() | ||
|
|
||
| flow.kickoff(user_message="أين طلبي؟", session_id=session_id) | ||
| flow.kickoff(user_message="وماذا عن الإرجاع؟", session_id=session_id) | ||
| flow.finalize_session_traces() | ||
| ``` | ||
|
|
||
| ## دورة حياة الجولة | ||
|
|
||
| كل `kickoff` مع `user_message` يشغّل: | ||
|
|
||
| 1. **`_configure_conversational_kickoff`** — دمج `session_id` / `user_message` في `inputs` وتطبيق `ConversationalConfig`. | ||
| 2. **استعادة الحالة** — عند وجود `inputs["id"]` و`@persist`. | ||
| 3. **`FlowStarted`** — في أول جولة للجلسة المؤجلة فقط. | ||
| 4. **`prepare_conversational_turn`** — إضافة رسالة المستخدم و`last_user_message` وتصنيف اختياري. | ||
| 5. **تنفيذ الرسم** — `@start` → `@router` → معالجات `@listen`. | ||
| 6. **نهاية التشغيل** — يُتخطى `flow_finished` والتتبع لكل جولة عند التأجيل؛ `Agent.kickoff()` / crews لا تغلق دفعة الأب. | ||
|
|
||
| استدعِ **`append_message("assistant", reply)`** في المعالجات. سطر المستخدم محفوظ عند kickoff — لا تُضفه مرة أخرى. | ||
|
|
||
| ## `ConversationalConfig` (افتراضيات على مستوى الصنف) | ||
|
|
||
| عيّن على صنف `Flow` كـ `conversational_config: ClassVar[ConversationalConfig | None]`. | ||
|
|
||
| | الحقل | الافتراضي | الغرض | | ||
| |-------|-----------|--------| | ||
| | `default_intents` | `None` | تسميات outcome للتصنيف التلقائي قبل kickoff | | ||
| | `intent_llm` | `None` | نموذج التصنيف (مطلوب عند وجود intents) | | ||
| | `interactive_prompt` | `"You: "` | مطالبة `kickoff(interactive=True)` | | ||
| | `interactive_timeout` | `None` | مهلة لكل سطر في الوضع التفاعلي | | ||
| | `exit_commands` | `exit`, `quit` | كلمات إنهاء الوضع التفاعلي | | ||
| | `defer_trace_finalization` | `True` | إبقاء دفعة trace واحدة مفتوحة بين الجولات | | ||
|
|
||
| يمكن التجاوز لكل kickoff عبر `intents=` و`intent_llm=`. | ||
|
|
||
| ## `ChatState` (شكل الحالة الموصى به للحفظ) | ||
|
|
||
| ```python | ||
| from crewai.flow import ChatState | ||
|
|
||
|
|
||
| class MyChatState(ChatState): | ||
| # موروث: id, messages, last_user_message, last_intent, session_ready | ||
| research_turn_count: int = 0 | ||
| custom_flag: bool = False | ||
| ``` | ||
|
|
||
| | الحقل | الدور | | ||
| |-------|------| | ||
| | `id` | UUID الجلسة (مثل `session_id` / `inputs["id"]`) | | ||
| | `messages` | قائمة `{role, content}` لسجل LLM | | ||
| | `last_user_message` | آخر سطر مستخدم في هذه الجولة | | ||
| | `last_intent` | تسمية المسار بعد التصنيف (إن وُجد) | | ||
| | `session_ready` | علم bootstrap لمرة واحدة | | ||
|
|
||
| `ConversationalInputs` هو `TypedDict` لـ `kickoff(inputs={...})`: `id`, `user_message`, `last_intent`. | ||
|
|
||
| ## API المحادثة على `Flow` | ||
|
|
||
| ### معاملات `kickoff` / `kickoff_async` | ||
|
|
||
| | المعامل | الغرض | | ||
| |---------|--------| | ||
| | `user_message` | نص هذه الجولة (أو `{"role": "user", "content": "..."}`) | | ||
| | `session_id` | UUID المحادثة → `inputs["id"]` / `state.id` | | ||
| | `intents` | تسميات outcome لـ `classify_intent` قبل kickoff | | ||
| | `intent_llm` | LLM للتصنيف (مطلوب مع `intents`) | | ||
| | `interactive` | حلقة CLI عبر `ask()` (للعروض المحلية فقط) | | ||
| | `interactive_prompt` | مطالبة الوضع التفاعلي | | ||
| | `interactive_timeout` | مهلة `ask()` لكل سطر | | ||
| | `exit_commands` | كلمات إنهاء الوضع التفاعلي | | ||
| | `inputs` | حقول حالة إضافية | | ||
| | `restore_from_state_id` | استنساخ من flow محفوظ آخر | | ||
|
|
||
| ### سمات المثيل | ||
|
|
||
| | السمة | الغرض | | ||
| |-------|--------| | ||
| | `conversational_config` | افتراضيات `ConversationalConfig` على مستوى الصنف | | ||
| | `defer_trace_finalization` | علم المثيل؛ يُضبط تلقائياً من config عند kickoff | | ||
| | `suppress_flow_events` | يخفي لوحات console؛ **التتبع يُسجّل** | | ||
| | `stream` | بث؛ مع `ChatSession.handle_turn(..., stream=True)` | | ||
|
|
||
| ### طرق وخصائص | ||
|
|
||
| | الاسم | الوصف | | ||
| |------|--------| | ||
| | `append_message(role, content, **extra)` | إضافة إلى `state.messages` | | ||
| | `conversation_messages` | سجل للقراءة فقط لاستدعاءات LLM | | ||
| | `classify_intent(text, outcomes, *, llm, context=None)` | تعيين outcome | | ||
| | `receive_user_message(text, *, outcomes=None, llm=None)` | إضافة رسالة مستخدم؛ `last_intent` اختياري | | ||
| | `finalize_session_traces()` | إصدار `flow_finished` المؤجل وإنهاء دفعة trace | | ||
| | `_should_defer_trace_finalization()` | هل يُؤجل إنهاء trace لكل جولة | | ||
| | `input_history` | سجل تدقيق مطالبات وردود `ask()` | | ||
|
|
||
| ### مساعدات الوحدة (`crewai.flow.conversation`) | ||
|
|
||
| | الدالة | الوصف | | ||
| |--------|--------| | ||
| | `normalize_kickoff_inputs(...)` | دمج kwargs المحادثة في `inputs` | | ||
| | `get_conversation_messages(flow)` | قراءة الرسائل من الحالة أو المخزن | | ||
| | `append_message(flow, ...)` | مثل طريقة المثيل | | ||
| | `prepare_conversational_turn(flow, ...)` | تهيئة الجولة (عادةً kickoff يستدعيها) | | ||
| | `receive_user_message(flow, ...)` | مثل طريقة المثيل | | ||
| | `set_state_field(flow, name, value)` | تعيين حقل dict أو Pydantic | | ||
| | `get_conversational_config(flow)` | قراءة `conversational_config` | | ||
| | `input_history_to_messages(entries)` | تحويل `input_history` لصيغة رسائل LLM | | ||
|
|
||
| ## أنماط توجيه النية | ||
|
|
||
| ### أ. تصنيف مسبق عبر `ConversationalConfig` (الأبسط) | ||
|
|
||
| عيّن `default_intents` و`intent_llm`. كل kickoff يصنّف قبل `@router`؛ اقرأ `self.state.last_intent` في `route()`. | ||
|
|
||
| ### ب. تصنيف داخل `@router` (مطالبات أغنى) | ||
|
|
||
| عيّن `default_intents=None` ليضيف kickoff الرسالة فقط. في `route()` استدعِ `classify_intent`: | ||
|
|
||
| ```python | ||
| @router(bootstrap) | ||
| def route(self): | ||
| intent = self.classify_intent( | ||
| self._routing_prompt(self.state.last_user_message), | ||
| ("GREETING", "ORDER", "RESEARCH", "GOODBYE"), | ||
| llm=self.conversational_config.intent_llm or "gpt-4o-mini", | ||
| ) | ||
| self.state.last_intent = intent | ||
| return intent | ||
| ``` | ||
|
|
||
| للبحث على الويب أو أدوات متعددة الخطوات استخدم **`@listen("RESEARCH")`** مع `Agent.kickoff()` وأدوات — وليس `LLM.call()` فقط. | ||
|
|
||
| ## عندما ينتهي الـ flow ويستمر المستخدم | ||
|
|
||
| `FlowFinished` يعني أن **تنفيذ الرسم هذا** اكتمل. تستمر المحادثة بـ `kickoff` آخر ونفس `session_id`. `@persist` يستعيد `messages` والأعلام والسياق. | ||
|
|
||
| **نمط الحفظ:** يُفضّل `@persist` على **خطوة نهائية واحدة** (مثل `finalize`) وليس على صنف `Flow` بالكامل. الحفظ على مستوى الصنف بعد كل method قد يفقد تحديثات المعالجات في نفس الجولة. | ||
|
|
||
| لا تستخدم `@human_feedback` لأسطر المتابعة في الدردشة إلا عند الحاجة لموافقة بشرية على مخرجات خطوة محددة. | ||
|
|
||
| ## التتبع عبر الجولات | ||
|
|
||
| مع `defer_trace_finalization=True` (افتراضي في `ConversationalConfig`): | ||
|
|
||
| - **دفعة trace واحدة** لجلسة الدردشة. | ||
| - **`flow_started`** في الجولة الأولى فقط؛ **`flow_finished`** مرة في `finalize_session_traces()`. | ||
| - **`kickoff` لكل جولة** لا يطبع "Trace batch finalized". | ||
| - **العمل المتداخل** (`Agent.kickoff()`, crews, Exa) يُلحق بدفعة **الأب**؛ flow داخلي من `AgentExecutor` لا يغلق دفعة الجلسة مبكراً. | ||
|
|
||
| ```python | ||
| try: | ||
| while True: | ||
| line = input("You: ").strip() | ||
| if not line: | ||
| break | ||
| flow.kickoff(user_message=line, session_id=session_id) | ||
| finally: | ||
| flow.finalize_session_traces() | ||
| ``` | ||
|
|
||
| `ChatSession.close()` يستدعي `finalize_session_traces()` عند التأجيل. | ||
|
|
||
| `suppress_flow_events=True` يخفي لوحات Rich فقط؛ أحداث trace والـ methods تُصدر. | ||
|
|
||
| ### `ConversationalFlow` عالي المستوى | ||
|
|
||
| يستخدم التجريد التجريبي عالي المستوى في `crewai.experimental.ConversationalFlow` دورة حياة tracing نفسها. عند تزيين الصنف بـ `@ConversationConfig(...)`، تكون قيمة `defer_trace_finalization` الافتراضية `True`، لذلك يبقي كل `handle_turn()` أثر الجلسة مفتوحاً بدلاً من تصديره فوراً. | ||
|
|
||
| إذا استدعيت `handle_turn()` مباشرةً، أنهِ الجلسة داخل كتلة `finally`: | ||
|
|
||
| ```python | ||
| from crewai.experimental import ConversationConfig, ConversationalFlow | ||
|
|
||
|
|
||
| @ConversationConfig() | ||
| class MyChatFlow(ConversationalFlow): | ||
| ... | ||
|
|
||
|
|
||
| flow = MyChatFlow() | ||
|
|
||
| try: | ||
| flow.handle_turn("مرحباً") | ||
| flow.handle_turn("أخبرني المزيد") | ||
| finally: | ||
| flow.finalize_session_traces() | ||
| ``` | ||
|
|
||
| إذا غلّفت الـ flow باستخدام `ChatSession`، فاستدعِ `session.close()` بدلاً من ذلك؛ فهو ينهي traces المؤجلة نيابةً عنك. بدون أحد هذين الاستدعاءين، قد تبقى دفعة trace مفتوحة وقد لا ترى trace مُصدَّراً لآخر محادثة. | ||
|
|
||
| ## `ChatSession` (WebSocket / SSE) | ||
|
|
||
| غلاف `kickoff` وجسر أحداث اختياري للواجهات. | ||
|
|
||
| ```python | ||
| from crewai.flow import ChatMessage, ChatSession | ||
|
|
||
|
|
||
| def on_event(msg: ChatMessage): | ||
| print(msg.type, msg.payload) | ||
|
|
||
|
|
||
| session = ChatSession( | ||
| flow, | ||
| session_id="channel-1", | ||
| intents=["order", "help"], | ||
| intent_llm="gpt-4o-mini", | ||
| on_event=on_event, | ||
| ) | ||
|
|
||
| turn = session.handle_turn("مرحباً") | ||
| print(turn.output, turn.intent, len(turn.messages)) | ||
|
|
||
| for msg in session.iter_turn_stream("أخبرني المزيد"): | ||
| print(msg.type, msg.payload) | ||
|
|
||
| session.close() # finalize_session_traces عند التأجيل | ||
| ``` | ||
|
|
||
| | النوع | الغرض | | ||
| |------|--------| | ||
| | `ChatSession` | جلسة؛ `handle_turn`, `iter_turn_stream`, `close` | | ||
| | `TurnResult` | `session_id`, `output`, `intent`, `messages`, `streaming` اختياري | | ||
| | `ChatMessage` | صيغة wire: `type`, `session_id`, `payload`, `seq` | | ||
| | `ConversationEventBridge` | bus → `ChatMessage` | | ||
|
|
||
| قيم `ChatMessage.type`: `user_message`, `assistant_delta`, `assistant_done`, `turn_started`, `turn_finished`, `error`, `tool_started`, `tool_finished`. | ||
|
|
||
| `stamp_conversation_fingerprint(event, session_id)` — يضبط `conversation_id` للموزّعات الخارجية. | ||
|
|
||
| ## `QueueInputProvider` (`ask` حاجز عبر WebSocket) | ||
|
|
||
| عندما يستدعي flow `ask()` داخل method وتصل الرسائل عبر socket: | ||
|
|
||
| ```python | ||
| from crewai.flow import Flow, QueueInputProvider | ||
|
|
||
| provider = QueueInputProvider() | ||
| flow = MyFlow(input_provider=provider) | ||
|
|
||
| provider.push(session_id, user_text) | ||
|
|
||
| reply = flow.ask("You: ", metadata={"session_id": session_id}) | ||
| provider.close_session(session_id) | ||
| ``` | ||
|
|
||
| ## البث | ||
|
|
||
| عيّن `stream = True` على صنف `Flow`. استخدم `kickoff(...)` أو `ChatSession.handle_turn(..., stream=True)` لأحداث مثل `assistant_delta`. | ||
|
|
||
| ## الاستيراد | ||
|
|
||
| ```python | ||
| from crewai.flow import ( | ||
| ChatMessage, | ||
| ChatSession, | ||
| ChatState, | ||
| ConversationalConfig, | ||
| ConversationalInputs, | ||
| ConversationEventBridge, | ||
| Flow, | ||
| QueueInputProvider, | ||
| TurnResult, | ||
| listen, | ||
| persist, | ||
| router, | ||
| start, | ||
| ) | ||
| ``` | ||
|
|
||
| ## مراجع | ||
|
|
||
| - [إتقان إدارة حالة Flow](/ar/guides/flows/mastering-flow-state) | ||
| - [أنشئ أول Flow](/ar/guides/flows/first-flow) | ||
| - Demo: `lib/crewai/runner_conversational_flow_simple.py` — REPL بسيط مع `RESEARCH` ووكيل Exa | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.