BF-Realtime / realtime_relay.py
SamiKoen
Prompt: tool cagrisi zorunlu β€” pronoun (alti, yedi) ile ilgili urun cagrilsin
ee0fdf9
Raw
History Blame Contribute Delete
17.4 kB
"""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