app2 / board /engine.py
Dooratre's picture
Upload 318 files
c37398c verified
Raw
History Blame Contribute Delete
25.9 kB
"""
Board AI Engine - STREAMING VERSION.
Updated: per-user locks (gevent), zero global blocking.
"""
import re
import json
import time
import gevent
from gevent.lock import RLock
from config import GPT_URL, MAX_CHAT_HISTORY, GPT_TIMEOUT
from http_pool import gpt_session
from board.tts import TTSEngine
from board.icons import IconResolver
from board.pages import resolve_page_tags
from subjects.loader import subject_loader
from json_processor import BoardProcessor
class BoardEngine:
"""
Board AI Engine with streaming segment delivery.
Per-user locks: User A never blocks User B.
"""
def __init__(self):
self.gpt_url = GPT_URL
self.board_processor = BoardProcessor()
self.tts_engine = TTSEngine()
self.icon_resolver = IconResolver()
# Per-user state with per-user locks
self._user_sessions = {}
self._user_locks = {}
self._master_lock = RLock() # Only for creating new user entries
# Cleanup inactive user locks every 30 minutes
gevent.spawn(self._periodic_lock_cleanup)
print("✅ BoardEngine initialized (per-user locks + gevent)")
# ─── Per-user lock management ───
def _get_user_lock(self, username):
"""Get or create a lock for a specific user. Only the master lock is held briefly."""
with self._master_lock:
if username not in self._user_locks:
self._user_locks[username] = {
'lock': RLock(),
'last_access': time.time()
}
self._user_locks[username]['last_access'] = time.time()
return self._user_locks[username]['lock']
def _periodic_lock_cleanup(self):
"""Remove locks for users inactive for 30+ minutes."""
while True:
gevent.sleep(1800) # 30 minutes
try:
now = time.time()
cutoff = now - 1800
with self._master_lock:
stale = [u for u, info in self._user_locks.items()
if info['last_access'] < cutoff and u not in self._user_sessions]
for u in stale:
del self._user_locks[u]
if stale:
print(f" 🗑️ Cleaned {len(stale)} stale user locks")
except Exception:
pass
# ─── User session management (per-user locked) ───
def _get_user_session(self, username):
lock = self._get_user_lock(username)
with lock:
if username not in self._user_sessions:
self._user_sessions[username] = {
"subject_id": None,
"conversation_history": [],
"last_sequence": [],
}
return self._user_sessions[username]
def set_subject(self, username, subject_id):
lock = self._get_user_lock(username)
with lock:
if username not in self._user_sessions:
self._user_sessions[username] = {
"subject_id": None,
"conversation_history": [],
"last_sequence": [],
}
us = self._user_sessions[username]
if us["subject_id"] and us["subject_id"] != subject_id:
us["conversation_history"] = []
us["last_sequence"] = []
print(f" 🔄 Board subject switched: {us['subject_id']}{subject_id} for {username}")
us["subject_id"] = subject_id
print(f" 📋 Board subject set: {subject_id} for {username}")
def get_subject(self, username):
lock = self._get_user_lock(username)
with lock:
us = self._user_sessions.get(username, {})
return us.get("subject_id")
# ─── GPT call (uses connection pool) ───
def _call_gpt5(self, user_message, system_prompt, temperature=0.7, max_tokens=4000):
payload = {
"user_input": user_message,
"chat_history": [
{"role": "system", "content": system_prompt}
],
"temperature": temperature,
"top_p": 0.95,
"max_completion_tokens": max_tokens
}
try:
response = gpt_session.post(self.gpt_url, json=payload, timeout=GPT_TIMEOUT)
response.raise_for_status()
return response.json().get("assistant_response", "")
except Exception as e:
print(f" ❌ GPT error: {e}")
return None
# ─── Chat history formatting ───
def _format_chat_history(self, username):
us = self._get_user_session(username)
history = us.get("conversation_history", [])
if not history:
return ""
recent = history[-10:]
parts = []
for msg in recent:
if msg["role"] == "user":
parts.append(f"الطالب: {msg['content']}")
elif msg["role"] == "assistant":
parts.append(f"المدرس: {msg['content']}")
return (
"\n\n═══ سجل المحادثة السابقة ═══\n"
+ "\n".join(parts)
+ "\n═══ نهاية السجل ═══"
)
# ─── Step 1: Route message ───
def _route_message(self, user_message, username):
us = self._get_user_session(username)
subject_id = us["subject_id"]
if not subject_id:
return "main.txt"
subject_data = subject_loader.load(subject_id)
if not subject_data:
return "main.txt"
structure = subject_data.get("structure.txt", "")
p_files = subject_data.get("_p_files", [])
chat_history_text = self._format_chat_history(username)
p_files_desc = "\n".join([f"- {f}: الفصل {i+1}" for i, f in enumerate(p_files)])
routing_prompt = f"""أنت نظام توجيه ذكي لمساعد تعليمي.
مهمتك: تحليل رسالة الطالب واختيار الملف المناسب للرد.
الملفات المتاحة:
- main.txt: للتحيات، الأسئلة العامة، أي شيء لا يتعلق بفصل محدد
{p_files_desc}
فهرس الكتاب (للمساعدة في التوجيه):
{structure}
{chat_history_text}
تعليمات:
1. إذا كانت الرسالة تحية أو سؤال عام → main.txt
2. إذا كانت تسأل عن موضوع في فصل محدد → الملف المناسب
3. استخدم الفهرس لتحديد الفصل الصحيح
4. مهم جداً: إذا قال الطالب "اشرح أكثر" أو "وضح" أو أي طلب متابعة، ارجع لسجل المحادثة لتعرف الموضوع الحالي واختر نفس الملف
5. أجب فقط باسم الملف (مثال: p1.txt أو main.txt) بدون أي كلام إضافي
رسالة الطالب: {user_message}
الملف المناسب:"""
chosen = self._call_gpt5(
user_message, routing_prompt,
temperature=0.2, max_tokens=50
)
if not chosen:
return "main.txt"
chosen = chosen.strip().lower()
valid_files = ["main.txt"] + p_files
for v in valid_files:
if v in chosen:
return v
return "main.txt"
# ─── Step 2: Generate XML response ───
def _generate_xml_response(self, user_message, chosen_file, username):
us = self._get_user_session(username)
subject_id = us["subject_id"]
if not subject_id:
return None
subject_data = subject_loader.load(subject_id)
if not subject_data:
return None
file_content = subject_data.get(chosen_file, "")
chat_history_text = self._format_chat_history(username)
system_prompt = f"""انتي مدرسة خبيرة ومحترفة في هذه المادة. تشرح على سبورة تفاعلية رقمية.
═══ صيغة الرد ═══
يجب أن تردي بصيغة XML خاصة تحتوي على:
1. <board>عناصر السبورة</board> - ما سيُرسم/يُضاف على السبورة أولاً
2. <voice>نص الكلام</voice> - النص الذي سيُقرأ بصوت عالٍ للطالب بعد رسم العناصر (عربي طبيعي)
═══ عناصر السبورة المتاحة (داخل <board>) ═══
• <note>محتوى الملاحظة</note>
• <text>نص مباشر على السبورة بدون خلفية</text>
• <shape type="نوع_الشكل"/>
الأنواع: circle, triangle, star, arrow-right, arrow-left, arrow-up, arrow-down,
rectangle, diamond, hexagon, square, oval, arrow-double-h, checkmark, cross,
heart, cloud, lightning, speech, process, decision
• <svg>كلمة_بحث_بالإنجليزية</svg>
سيتم البحث عن أيقونة مرسومة يدوياً (مثل: ball, car, force, spring, weight, rope, pulley)
• <page>رقم_الصفحة</page>
لعرض صفحة محددة من الكتاب كصورة على السبورة
مثال: <page>12</page> لعرض الصفحة 12 من الكتاب
استخدمها عندما تحتاج تعرض للطالب صفحة معينة من الكتاب
═══ قواعد مهمة جداً ═══
1. اشرح خطوة بخطوة: ابدأ بـ <board> ثم <voice> ثم <board> ثم <voice> وهكذا
2. اجعل الشرح متدرجاً كأنك تشرح على سبورة حقيقية أمام الطلاب
3. <board> = ما يظهر على السبورة أولاً (ملاحظات، نصوص، أشكال، صور، صفحات الكتاب)
4. <voice> = الكلام المسموع بعد رسم العناصر (طبيعي، ودود، واضح، يشرح ما تم رسمه)
5. لا تضع كل شيء دفعة واحدة - اجعله تسلسلياً
7. <svg> فقط بكلمات إنجليزية بسيطة ومعبرة
8. السبورة تعمل بنظام الإضافة - العناصر السابقة تبقى
9. استخدم <text> للعناوين والمعادلات المهمة (بدون خلفية)
10. استخدم <note> للتوضيحات والملاحظات (مع خلفية ملونة)
11. لا تستخدم أكثر من 3-4 عناصر في كل <board>
12. اجعل النص في <voice> طبيعياً كأنك تتحدث مع طالب ويشرح ما تم رسمه على السبورة
13. ارسم أولاً ثم تكلم - هذا مهم جداً!
14. راجع سجل المحادثة السابقة لتعرف ما تم شرحه وتكمل من حيث توقفت - لا تكرر ما قلته سابقاً
15. استخدم <page> عندما تريد عرض صفحة من الكتاب - مثلاً إذا الطالب سأل عن تمرين أو شكل في صفحة معينة
when user talk about something not about the subject or something funny etc... you can actually answer without the board just VOICE and be funny smart perfect girl also:
when you explain something dont make all your explain on the NOTE make the note for important point use the TEXT direct on the board and the ICONS/SHAPES
═══ محتوى المادة ═══
{file_content}
{chat_history_text}
═══ الآن أجب على سؤال الطالب ═══
رسالة الطالب: {user_message}"""
response = self._call_gpt5(
user_message, system_prompt,
temperature=0.8, max_tokens=4000
)
return response
# ─── Resolve tags ───
def _resolve_svg_tags(self, xml_text):
if not xml_text:
return xml_text
return self.icon_resolver.resolve_all_in_xml(xml_text)
def _resolve_page_tags(self, xml_text, username):
if not xml_text:
return xml_text
us = self._get_user_session(username)
subject_id = us.get("subject_id")
if not subject_id:
return xml_text
base_url = subject_loader.get_pages_base_url(subject_id)
if not base_url:
print(f" ⚠️ No pages_base_url for subject {subject_id}")
return xml_text
return resolve_page_tags(xml_text, base_url)
# ─── Parse XML into raw segments ───
def _parse_xml_to_raw_segments(self, xml_response):
if not xml_response:
return []
pattern = r'<(voice|board)>(.*?)</\1>'
matches = list(re.finditer(pattern, xml_response, re.DOTALL))
if not matches:
cleaned = re.sub(r'<[^>]+>', '', xml_response).strip()
if cleaned:
return [{"boards": [], "voice": cleaned}]
return []
groups = []
current_boards = []
for match in matches:
tag_type = match.group(1)
content = match.group(2).strip()
if tag_type == "board":
current_boards.append(content)
elif tag_type == "voice":
cleaned_voice = re.sub(r'<[^>]+>', '', content).strip()
groups.append({
"boards": list(current_boards),
"voice": cleaned_voice
})
current_boards = []
if current_boards:
groups.append({
"boards": list(current_boards),
"voice": ""
})
return groups
# ─── Process a single board content into items ───
def _process_single_board(self, board_content, current_board_state):
existing_json_str = json.dumps(
current_board_state, ensure_ascii=False, indent=2
)
processor_input = (
f"BOARD NOW (make sure no X Y error):\n"
f"{existing_json_str}\n\n"
f"new board need to add :\n"
f"<board>{board_content}</board>"
)
print(f" 🔧 Sending to json_processor...")
print(f" Current board items: {len(current_board_state)}")
try:
json_text = self.board_processor.convert_xml_to_json(processor_input)
if json_text:
new_items = json.loads(json_text)
if isinstance(new_items, list) and new_items:
added_items = []
existing_ids = set()
for item in current_board_state:
item_key = json.dumps(item, sort_keys=True, ensure_ascii=False)
existing_ids.add(item_key)
for item in new_items:
item_key = json.dumps(item, sort_keys=True, ensure_ascii=False)
if item_key not in existing_ids:
added_items.append(item)
current_board_state.extend(added_items)
print(f" ✅ json_processor: {len(new_items)} total, {len(added_items)} new")
return added_items, current_board_state
elif isinstance(new_items, dict):
current_board_state.append(new_items)
print(f" ✅ json_processor: 1 item")
return [new_items], current_board_state
else:
print(f" ⚠️ json_processor: unexpected format")
return [], current_board_state
except json.JSONDecodeError as e:
print(f" ❌ json_processor invalid JSON: {e}")
return [], current_board_state
except Exception as e:
print(f" ❌ json_processor error: {e}")
return [], current_board_state
return [], current_board_state
# ─── STREAMING PIPELINE (generator) ───
def process_message_stream(self, user_message, username, frontend_board_state=None):
"""
Stream board response. NO global lock - only per-user lock for session reads.
The heavy work (GPT calls, TTS) runs WITHOUT holding any lock.
"""
# Read session data under per-user lock (fast)
us = self._get_user_session(username)
subject_id = us.get("subject_id")
print(f"\n{'═' * 60}")
print(f" 👤 Student ({username}): {user_message}")
print(f" 📚 Subject: {subject_id}")
print(f"{'═' * 60}")
if not subject_id:
yield json.dumps({
"type": "error",
"message": "لم يتم تحديد المادة للسبورة"
}, ensure_ascii=False)
return
if frontend_board_state is None:
frontend_board_state = []
current_board_state = list(frontend_board_state)
print(f" 📋 Board state from frontend: {len(current_board_state)} items")
# Step 1: Route (NO lock held - just GPT call)
print("\n 📍 Step 1: Routing message...")
chosen_file = self._route_message(user_message, username)
print(f" 📂 Chosen file: {chosen_file}")
# Step 2: Generate XML (NO lock held - just GPT call)
print(f" 🤖 Step 2: Generating XML response...")
xml_response = self._generate_xml_response(user_message, chosen_file, username)
if not xml_response:
print(" ❌ Failed to generate response")
yield json.dumps({
"type": "error",
"message": "عذراً، حدث خطأ في النظام. حاول مرة أخرى."
}, ensure_ascii=False)
return
print(f" 📝 XML response: {len(xml_response)} chars")
# Step 3: Resolve <page> tags (NO lock)
print(" 📄 Step 3: Resolving <page> tags...")
xml_response = self._resolve_page_tags(xml_response, username)
# Step 4: Resolve <svg> tags (NO lock)
print(" 🎨 Step 4: Resolving <svg> tags...")
xml_response = self._resolve_svg_tags(xml_response)
# Step 5: Parse XML into segment groups (NO lock - pure CPU, fast)
print(" 🔧 Step 5: Parsing XML into segment groups...")
raw_groups = self._parse_xml_to_raw_segments(xml_response)
total_groups = len(raw_groups)
print(f" 📊 Found {total_groups} segment groups to stream")
if total_groups == 0:
yield json.dumps({
"type": "error",
"message": "لم يتم توليد محتوى للسبورة."
}, ensure_ascii=False)
return
# Step 6: Process and yield each group (NO lock during heavy work)
all_voice_texts = []
all_segments_for_replay = []
for idx, group in enumerate(raw_groups):
print(f"\n ── Segment {idx + 1}/{total_groups} ──")
segment_board_items = []
for board_content in group["boards"]:
new_items, current_board_state = self._process_single_board(
board_content, current_board_state
)
segment_board_items.extend(new_items)
voice_text = group.get("voice", "")
audio_url = None
if voice_text:
all_voice_texts.append(voice_text)
text_preview = voice_text[:50]
print(f" 🎙️ Converting TTS: '{text_preview}...'")
audio_file = self.tts_engine.convert(voice_text)
if audio_file:
audio_url = f"/static/{audio_file}"
else:
print(f" ⚠️ TTS failed for this segment")
segment_data = {
"type": "segment",
"index": idx,
"total_estimate": total_groups,
"board_items": segment_board_items,
"voice_text": voice_text,
"audio_url": audio_url
}
all_segments_for_replay.append(segment_data)
print(f" ✅ Segment {idx + 1} ready: {len(segment_board_items)} board items, voice={'yes' if voice_text else 'no'}")
yield json.dumps(segment_data, ensure_ascii=False)
# Step 7: Save chat history (per-user lock, fast)
lock = self._get_user_lock(username)
with lock:
us = self._user_sessions.get(username, {})
if "conversation_history" not in us:
us["conversation_history"] = []
us["conversation_history"].append({
"role": "user",
"content": user_message
})
assistant_text = " ".join(all_voice_texts)
if assistant_text:
us["conversation_history"].append({
"role": "assistant",
"content": assistant_text
})
if len(us["conversation_history"]) > MAX_CHAT_HISTORY:
us["conversation_history"] = us["conversation_history"][-MAX_CHAT_HISTORY:]
us["last_sequence"] = all_segments_for_replay
# Trigger cleanup in background (non-blocking)
self.tts_engine.request_cleanup()
yield json.dumps({
"type": "done",
"board_state": current_board_state,
"chosen_file": chosen_file,
"total_segments": total_groups
}, ensure_ascii=False)
print(f"\n ✅ All {total_groups} segments streamed!")
print(f" Board items: {len(current_board_state)}")
print(f"{'═' * 60}\n")
# ─── NON-STREAMING (kept for compatibility) ───
def process_message(self, user_message, username, frontend_board_state=None):
all_segments = []
final_board_state = frontend_board_state or []
chosen_file = None
for chunk_str in self.process_message_stream(user_message, username, frontend_board_state):
chunk = json.loads(chunk_str)
if chunk["type"] == "segment":
if chunk.get("board_items"):
all_segments.append({
"type": "board_update",
"action": "add",
"items": chunk["board_items"]
})
if chunk.get("voice_text"):
all_segments.append({
"type": "voice",
"text": chunk["voice_text"],
"audio_url": chunk.get("audio_url")
})
elif chunk["type"] == "done":
final_board_state = chunk.get("board_state", [])
chosen_file = chunk.get("chosen_file")
elif chunk["type"] == "error":
return {
"success": False,
"error": chunk["message"],
"sequence": [{
"type": "voice",
"text": chunk["message"],
"audio_url": None
}],
"board_state": frontend_board_state or []
}
return {
"success": True,
"chosen_file": chosen_file,
"sequence": all_segments,
"board_state": final_board_state
}
# ─── Replay ───
def get_replay_sequence(self, username):
lock = self._get_user_lock(username)
with lock:
us = self._user_sessions.get(username, {})
last_seq = us.get("last_sequence", [])
if not last_seq:
return {
"success": False,
"error": "No previous response to replay",
"sequence": []
}
voice_only = []
for item in last_seq:
if item.get("voice_text"):
voice_only.append({
"type": "voice",
"text": item.get("voice_text", ""),
"audio_url": item.get("audio_url")
})
print(f" 🔄 Replay: {len(voice_only)} voice segments")
return {
"success": True,
"sequence": voice_only
}
# ─── Clear ───
def clear_board(self, username):
lock = self._get_user_lock(username)
with lock:
if username in self._user_sessions:
self._user_sessions[username]["last_sequence"] = []
print(f" 🗑️ Board cleared for {username}")
return {"success": True, "board_state": []}
def clear_chat_history(self, username):
lock = self._get_user_lock(username)
with lock:
if username in self._user_sessions:
self._user_sessions[username]["conversation_history"] = []
self._user_sessions[username]["last_sequence"] = []
print(f" 🗑️ Board chat history cleared for {username}")
return {"success": True}
def clear_user_session(self, username):
lock = self._get_user_lock(username)
with lock:
if username in self._user_sessions:
del self._user_sessions[username]
print(f" 🗑️ Full board session cleared for {username}")
return {"success": True}