"""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