Spaces:
Running
Running
| """OpenAI Realtime API <-> Browser WebSocket relay. | |
| Akis: | |
| 1. Browser baglanir | |
| 2. OpenAI Realtime'a baglanip session config gonderir | |
| 3. Iki yonlu mesaj relay'i: | |
| - Browser -> OpenAI: PCM audio + control | |
| - OpenAI -> Browser: audio chunks + transcript + tool calls + custom events | |
| 4. Tool call'lari yakala, thread'de calistir, sonucu OpenAI'ye geri gonder | |
| 5. Asistan transkripti + kullanici intent'iyle urun gorsellerini live update et | |
| """ | |
| from __future__ import annotations | |
| import asyncio | |
| import json | |
| import logging | |
| import websockets | |
| from websockets.asyncio.client import connect as ws_connect | |
| from fastapi import WebSocket, WebSocketDisconnect | |
| from browser_session import get_browser_session | |
| from config import OPENAI_API_KEY, REALTIME_MODEL, REALTIME_URL | |
| from product_matcher import ( | |
| extract_product_link_from_text, | |
| find_color_variant, | |
| find_main_local_substring, | |
| find_main_product_in_text, | |
| find_product_in_transcript, | |
| ) | |
| from product_index import get_index | |
| from prompts import get_active_prompt_content_only | |
| from tools import TOOLS, handle_tool_call_sync, strip_urls | |
| logger = logging.getLogger(__name__) | |
| VOICE_ADDON = ( | |
| "\n\nGORSEL DESTEK (onemli):\n" | |
| "- Musterinin ekraninin sag tarafinda urun gorseli OTOMATIK olarak gosteriliyor.\n" | |
| "- Bahsettigin urunun resmi sag tarafta belirir; renk sorulursa o varyantin gorseline gecer.\n" | |
| "- 'Gorsel paylasiyorum', 'resmini atiyorum', 'linki gonderiyorum' gibi cumleler kurma.\n" | |
| "- 'Sagda gorebilirsiniz' gibi kisa referanslar verebilirsin.\n" | |
| "\nTOOL CAGRISI ZORUNLU (en kritik kural):\n" | |
| "- Stok/mevcudiyet/depo/fiyat/renk/beden ile ilgili HERHANGI bir soruda check_warehouse_stock CAGIR.\n" | |
| "- 'var mi?', 'stokta mi?', 'kaΓ§ para?', 'fiyati ne?', 'hangi magazada?', 'X'da var mi?', " | |
| "'altΔ± var mi?', 'yedi var mi?', 'iki tane var mi?' β HEPSI tool gerektirir.\n" | |
| "- Pronoun'lar (altΔ±, yedi, beΕi, onun, bunun, sunun) ONCEKI urunun model numarasi olabilir. " | |
| "Ornek: 'Marlin 5 var mi?' -> sonra 'AltΔ± var mi?' = 'Marlin 6 var mi?'. Bagle, tool cagir.\n" | |
| "- Tool olmadan stok/fiyat/depo cevabi VERME. 'Bilmiyorum' veya 'kontrol edeyim' deme β " | |
| "hemen tool cagir.\n" | |
| "- Yeni urun adi geΓ§tiyse her zaman show_product VEYA check_warehouse_stock cagir.\n" | |
| "\nIKI TOOL VAR β DOGRU OLANI SEC (kritik):\n" | |
| "- show_product (HAFIF, varsayilan): urun BAHSEDILDIGINDE sayfayi sag monitorde acar. " | |
| "Fiyat, gorsel, aciklama sayfada gorunur. XML/stok cekmez. HIZLI VE UCUZ.\n" | |
| "- check_warehouse_stock (AGIR, sadece gerektiginde): SADECE musteri stok/depo durumu " | |
| "sorarsa kullan. Ornek: 'var mi?', 'stokta mi?', 'Caddebostan'da var mi?', " | |
| "'hangi magazada bulurum?'. Bu tool BizimHesap'a istek atar β gereksizse cagirma.\n" | |
| "- HER spesifik urun adi soyleyecegin ZAMAN β bisiklet, sele, kask, far, jersey, jant, " | |
| "gidon, pedal, lastik, ne olursa olsun β once show_product cagir (sayfa acilsin).\n" | |
| "- Ornekler:\n" | |
| " * Musteri 'Marlin 5 hakkinda bilgi' -> show_product('Marlin 5')\n" | |
| " * Musteri 'Madone SLR 9' -> show_product('Madone SLR 9')\n" | |
| " * Musteri 'Marlin 5 var mi?' -> check_warehouse_stock('Marlin 5 var mi?')\n" | |
| " * Musteri 'Caddebostan'da Madone var mi?' -> check_warehouse_stock(...)\n" | |
| "- Onceki cevapta bahsedilmis olsa bile, YENI bir model ismi soyleyince TEKRAR cagir.\n" | |
| "- Tek istisna: marka/kategori adi (Trek, Bontrager, MTB, sele, kask gibi jenerik kelime).\n" | |
| "- TAM URUN ADI VER (kategori dahil):\n" | |
| " * 'Comp' degil 'Bontrager Comp sele'\n" | |
| " * 'Pro' degil 'Bontrager Aeolus Pro sele'\n" | |
| " * 'XR4' degil 'Bontrager XR4 Team Issue lastik'\n" | |
| " Marka + model + KATEGORI (sele, jant, gidon, pedal, lastik, kask...).\n" | |
| "\nSESLI SOHBET KURALLARI (cok onemli):\n" | |
| "- Hedef: TAM 2 cumle. Asla 2'den fazla cumle KURMA. 3. cumle YASAK.\n" | |
| "- Bilgi 1 cumlede sigarsa 1 cumle yeterli. Sigmazsa 2 cumle. Daha fazlasi yok.\n" | |
| "- IFADELERINI CESITLENDIR. Ayni kalipta cumle kurma β her cevapta farkli kelime ve yapi sec.\n" | |
| "- Sadece SORULANA dogrudan cevap ver. Ekstra oneri/yonlendirme/tesvik EKLEME.\n" | |
| "- KESINLIKLE YASAK kapanis/dolgu cumleleri (cevabini bunlarla UZATMA, son sozun bunlar OLMASIN):\n" | |
| " * 'magazalarimizda inceleyebilirsiniz', 'magazada gorebilirsiniz', 'magazada deneyebilirsiniz'\n" | |
| " * 'yardimci olabilir miyim', 'baska sorunuz var mi', 'baska bir sey...'\n" | |
| " * 'detayli bilgi icin', 'isterseniz', 'dilerseniz', 'arzu ederseniz'\n" | |
| " * 'bizi tercih ettiginiz icin', 'iyi gunler', 'kolay gelsin'\n" | |
| " Cevabi bilgi ile bitir, ek cumle EKLEME. Sustugun an iletisim biter.\n" | |
| "- Markdown, * veya emoji KULLANMA.\n" | |
| "- URL/link/web adresi ASLA SOYLEME. 'www', 'trekbisiklet.com', 'https' gibi adresleri " | |
| "hicbir sekilde sesli okuma. 'Linki paylasiyorum', 'sitede goreceksiniz' gibi link " | |
| "referanslari da YASAK.\n" | |
| "- HER ZAMAN 'siz' ile hitap et, soru ile bitirme.\n" | |
| "- Stok sorulari icin check_warehouse_stock cagir, sonucu OZUN tek cumlede ver.\n" | |
| "\n" | |
| "STOK ADEDI KURALI (kritik):\n" | |
| "- Stok ADEDINI / SAYISINI ASLA SOYLEME.\n" | |
| "- Sadece 'mevcut', 'stokta var', 'bulunmuyor' veya 'tukenmis' gibi durum bildir.\n" | |
| "- Adet sorulsa bile 'detayli adet bilgisi icin magazayla teyit' deyip gec.\n" | |
| "\n" | |
| "STOK SORGUSU AKISI:\n" | |
| "- check_warehouse_stock cagrisi 2-3 sn surebilir. Once KISA bir bekleme cumlesi soyle, " | |
| "SONRA fonksiyonu cagir. show_product ANINDA dondugu icin bekleme cumlesi gerekmez.\n" | |
| "- Bekleme cumlesi ornekleri (cesitlendir, hep ayni demeyi): " | |
| "'Bir saniye, bakiyorum.' / 'Hemen kontrol ediyorum.' / 'Bakiyorum efendim.' / " | |
| "'Stoga bakiyorum.' / 'Birsaniyenize.'\n" | |
| "- Sonuc gelince direkt cevaba gec; 'kontrol ettim' gibi tekrar cumlesi YOK.\n" | |
| ) | |
| def build_session_instructions() -> str: | |
| try: | |
| base = get_active_prompt_content_only() | |
| if isinstance(base, list): | |
| base = "\n\n".join(str(p) for p in base) | |
| except Exception: | |
| logger.exception("Prompt yuklenemedi") | |
| base = "Trek Bisiklet uzmani bir satis temsilcisisin." | |
| return base + VOICE_ADDON | |
| def _session_update_payload() -> dict: | |
| return { | |
| "type": "session.update", | |
| "session": { | |
| "type": "realtime", | |
| "model": REALTIME_MODEL, | |
| "instructions": build_session_instructions(), | |
| "output_modalities": ["audio"], | |
| "audio": { | |
| "input": { | |
| "format": {"type": "audio/pcm", "rate": 24000}, | |
| "transcription": {"model": "gpt-4o-mini-transcribe", "language": "tr"}, | |
| "turn_detection": { | |
| "type": "server_vad", | |
| "threshold": 0.7, # Daha katΔ± (0.5 -> 0.7) β asistan sesi user gibi algΔ±lanmasΔ±n | |
| "prefix_padding_ms": 300, | |
| "silence_duration_ms": 1200, # Sessizlik daha uzun (700 -> 1200ms) | |
| "interrupt_response": False, | |
| "create_response": True, | |
| }, | |
| }, | |
| "output": { | |
| "format": {"type": "audio/pcm", "rate": 24000}, | |
| }, | |
| }, | |
| "tools": TOOLS, | |
| "tool_choice": "auto", | |
| }, | |
| } | |
| async def realtime_relay(client_ws: WebSocket): | |
| """FastAPI WebSocket endpoint handler.""" | |
| await client_ws.accept() | |
| if not OPENAI_API_KEY: | |
| await client_ws.send_text(json.dumps({ | |
| "type": "error", | |
| "error": {"message": "OPENAI_API_KEY tanimli degil."}, | |
| })) | |
| await client_ws.close() | |
| return | |
| # Pre-warm index (varsa cache'den hizli, yoksa fetch) | |
| try: | |
| await asyncio.to_thread(get_index().ensure) | |
| except Exception: | |
| logger.exception("index pre-warm hatasi") | |
| headers = {"Authorization": f"Bearer {OPENAI_API_KEY}"} | |
| # Session state | |
| state = { | |
| "text": "", # asistan response transkripti | |
| "tool_link": None, # bu turn'de tool ile gosterilen urun linki | |
| "current_main_link": None, # ekrandaki ana urun linki | |
| "last_shown_in_response": None, # bu response'da en son gosterilen | |
| "last_check_len": 0, # live detection throttle | |
| } | |
| try: | |
| async with ws_connect(REALTIME_URL, additional_headers=headers) as openai_ws: | |
| logger.info("OpenAI Realtime baglantisi kuruldu") | |
| await openai_ws.send(json.dumps(_session_update_payload())) | |
| async def client_to_openai(): | |
| try: | |
| while True: | |
| msg = await client_ws.receive_text() | |
| await openai_ws.send(msg) | |
| except WebSocketDisconnect: | |
| logger.info("Client disconnected") | |
| except Exception as e: | |
| logger.error(f"client_to_openai error: {e}") | |
| async def openai_to_client(): | |
| try: | |
| async for raw in openai_ws: | |
| try: | |
| data = json.loads(raw) | |
| except Exception: | |
| await client_ws.send_text(raw) | |
| continue | |
| await _handle_event(data, raw, openai_ws, client_ws, state) | |
| except websockets.exceptions.ConnectionClosed: | |
| logger.info("OpenAI WebSocket kapandi") | |
| except Exception as e: | |
| logger.error(f"openai_to_client error: {e}") | |
| await asyncio.gather(client_to_openai(), openai_to_client()) | |
| except Exception: | |
| logger.exception("Realtime relay hatasi") | |
| try: | |
| await client_ws.send_text(json.dumps({ | |
| "type": "error", | |
| "error": {"message": "Baglanti hatasi"}, | |
| })) | |
| except Exception: | |
| pass | |
| finally: | |
| try: | |
| await client_ws.close() | |
| except Exception: | |
| pass | |
| async def _handle_event(data: dict, raw: str, openai_ws, client_ws: WebSocket, state: dict): | |
| evt = data.get("type", "") | |
| # ---------- Important event logging ---------- | |
| if evt in ("session.created", "session.updated", "input_audio_buffer.speech_started", | |
| "input_audio_buffer.speech_stopped", "input_audio_buffer.committed", | |
| "response.created"): | |
| logger.info(f"[Realtime] {evt}") | |
| elif evt == "error": | |
| logger.error(f"[Realtime] ERROR: {json.dumps(data)[:400]}") | |
| elif evt == "response.done": | |
| status = data.get("response", {}).get("status") | |
| details = data.get("response", {}).get("status_details") | |
| logger.info(f"[Realtime] response.done status={status} details={details}") | |
| # ---------- State resets ---------- | |
| # Asistan yeni response'a basliyor β assistant transcript ve tool flag reset | |
| if evt == "response.created": | |
| state["text"] = "" | |
| state["last_shown_in_response"] = None | |
| state["last_check_len"] = 0 | |
| state["pending_response_create"] = False | |
| # Kullanici yeni soru sormaya basladi β tool_link gecersiz | |
| if evt == "input_audio_buffer.speech_started": | |
| state["tool_link"] = None | |
| # ---------- HIZLI YOL: kullanici transcription tamamlandiginda ---------- | |
| # LLM'in tool secmesini beklemeden, transcript'te urun adi varsa local | |
| # substring matcher ile aninda browser'i navigate et. Sayfa kullanici | |
| # sustuktan ~500ms sonra acilir, asistan cevabi sonra gelir. | |
| if evt == "conversation.item.input_audio_transcription.completed": | |
| transcript = (data.get("transcript") or "").strip() | |
| if transcript: | |
| logger.info(f"[user-transcript] {transcript[:100]!r}") | |
| # Catalog adi transcript'in icinde gecirken aninda bul | |
| try: | |
| p = find_product_in_transcript(transcript) | |
| if p and p.get("link"): | |
| link = p["link"] | |
| if link != state.get("current_main_link"): | |
| logger.info(f"[fast-nav] {p['name']} -> {link}") | |
| state["current_main_link"] = link | |
| state["last_shown_in_response"] = link | |
| await _send(client_ws, {"type": "product.show", "product": p}) | |
| asyncio.create_task(get_browser_session().navigate(link, fallback_query=p.get("name"))) | |
| except Exception: | |
| logger.exception("fast-nav hatasi") | |
| # ---------- Asistan transkripti β sadece renk override kontrolu icin ---------- | |
| # GPT'ye guveniyoruz: ana urun gosterimi YALNIZCA tool sonucundan gelir. | |
| # Live transkript-bazli tahmin kapali (yanlis urun gosterimine yol aciyordu). | |
| if evt in ("response.audio_transcript.delta", | |
| "response.output_audio_transcript.delta", | |
| "response.output_text.delta"): | |
| d = data.get("delta", "") | |
| if isinstance(d, str): | |
| state["text"] += d | |
| # ---------- response.done β renk override + bekleyen response.create gonder ---------- | |
| if evt == "response.done": | |
| text = state["text"] | |
| active_main = state["current_main_link"] | |
| if text and active_main: | |
| color_v = find_color_variant(active_main, text) | |
| if color_v and color_v.get("link"): | |
| logger.info(f"[product/color-override] {color_v['name']}") | |
| await _send(client_ws, { | |
| "type": "product.show", | |
| "product": color_v, | |
| }) | |
| # Bu turn'de tool cagrisi yapildiysa SADECE BIR KEZ response.create gonder | |
| # (cogu tool cagrisi varsa bile sonuclari beraber donar, tek sozlu cevap olur). | |
| if state.get("pending_response_create"): | |
| state["pending_response_create"] = False | |
| await openai_ws.send(json.dumps({"type": "response.create"})) | |
| # ---------- Tool call yakalama ---------- | |
| if evt == "response.function_call_arguments.done": | |
| await _handle_tool_call(data, openai_ws, client_ws, state) | |
| # Tool call mesajini client'a forward etmeye gerek yok ama geri gondermek | |
| # zarar vermez β relay kuralina uy. | |
| # ---------- Tum mesajlari client'a forward et ---------- | |
| try: | |
| await client_ws.send_text(raw) | |
| except Exception: | |
| pass | |
| async def _handle_tool_call(data: dict, openai_ws, client_ws: WebSocket, state: dict): | |
| call_id = data.get("call_id") | |
| fn_name = data.get("name") | |
| try: | |
| args = json.loads(data.get("arguments", "{}")) | |
| except Exception: | |
| args = {} | |
| logger.info(f"Tool call: {fn_name}({args})") | |
| # Tool'u thread'de calistir β event loop bloklanmasin | |
| result = await asyncio.to_thread(handle_tool_call_sync, fn_name, args) | |
| # Urun resmi tespit β SADECE tool sonucundaki productLink'i kullan. | |
| # GPT'ye guven: hangi urun konusulduysa o linki dondurur, biz tahmin yurutmeyiz. | |
| if fn_name in ("show_product", "check_warehouse_stock", "get_warehouse_stock"): | |
| idx = get_index() | |
| product = None | |
| product_link = extract_product_link_from_text(result) | |
| if product_link: | |
| product = idx.find_by_link(product_link) | |
| # Fresh-upgrade: tool eski/obsolete urun bulduysa (link 404 olabilir), | |
| # ayni isimden fresh kardes varsa onu kullan | |
| if product: | |
| fresh = find_main_local_substring(product.get("name") or "") | |
| if fresh and fresh.get("link") and fresh.get("link") != product.get("link"): | |
| logger.info(f"[fresh-upgrade] {product.get('name')} -> {fresh.get('name')}") | |
| product = fresh | |
| if product and product.get("link"): | |
| logger.info(f"[product/tool] {product['name']} ({product.get('link')})") | |
| state["tool_link"] = product.get("link") or product_link | |
| if product.get("link"): | |
| state["current_main_link"] = product["link"] | |
| state["last_shown_in_response"] = product["link"] | |
| await _send(client_ws, { | |
| "type": "product.show", | |
| "product": product, | |
| }) | |
| # Sag monitor: headless browser'i bu sayfaya navige et | |
| # 404 ise urun adiyla Trek arama sayfasina fallback | |
| target_url = product.get("link") or product_link | |
| fallback_q = product.get("name") | |
| if target_url: | |
| asyncio.create_task(get_browser_session().navigate(target_url, fallback_query=fallback_q)) | |
| # Modelin URL okumasi engellensin β strip | |
| result = strip_urls(result) | |
| # Tool sonucunu modele geri gonder. response.create'i SADECE response.done'da | |
| # bir kez tetikle (paralel tool cagrilarinda duplicate audio uretiyordu). | |
| await openai_ws.send(json.dumps({ | |
| "type": "conversation.item.create", | |
| "item": { | |
| "type": "function_call_output", | |
| "call_id": call_id, | |
| "output": result, | |
| }, | |
| })) | |
| state["pending_response_create"] = True | |
| async def _send(client_ws: WebSocket, payload: dict): | |
| try: | |
| await client_ws.send_text(json.dumps(payload)) | |
| except Exception: | |
| pass | |