1414"""
1515
1616import json
17+ import logging
18+ import time
1719from dataclasses import dataclass , field
18- from typing import Any , Iterable
20+ from typing import Any , Callable , Iterable , TypeVar
1921
2022import httpx
2123
2224from gooddata_eval .core .models import ChatResult , DatasetItem
2325
26+ _log = logging .getLogger (__name__ )
27+
2428SSE_DATA_PREFIX = "data: "
2529
30+ _RETRYABLE_STATUS_CODES : frozenset [int ] = frozenset ({429 , 502 , 503 , 504 })
31+ _METADATA_SYNC_MARKER = "METADATA_SYNC_IN_PROGRESS"
32+
33+
34+ class ChatError (RuntimeError ):
35+ """Non-retryable error reported by the chat SSE stream."""
36+
37+ def __init__ (self , message : str , * , status_code : int | None = None , detail : str | None = None ) -> None :
38+ super ().__init__ (message )
39+ self .status_code = status_code
40+ self .detail = detail
41+
42+
43+ class TransientChatError (ChatError ):
44+ """Retryable transient error: gen-ai temporarily unavailable or still syncing metadata."""
45+
46+
47+ _MAX_RETRIES = 5
48+ _INITIAL_BACKOFF_S = 5.0
49+ _BACKOFF_FACTOR = 2.0
50+ _MAX_BACKOFF_S = 60.0
51+
52+ T = TypeVar ("T" )
53+
54+
55+ def _is_retryable_exc (exc : Exception ) -> bool :
56+ if isinstance (exc , TransientChatError ):
57+ return True
58+ if isinstance (exc , httpx .HTTPStatusError ):
59+ return exc .response .status_code in _RETRYABLE_STATUS_CODES
60+ return False
61+
62+
63+ def _retry_transient (operation : Callable [[], T ], * , is_retryable : Callable [[Exception ], bool ]) -> T :
64+ """Run ``operation``; retry retryable failures with bounded exponential backoff."""
65+ delay = _INITIAL_BACKOFF_S
66+ for attempt in range (_MAX_RETRIES + 1 ): # 0..N => N retries + 1 initial attempt
67+ try :
68+ return operation ()
69+ except Exception as exc :
70+ if attempt == _MAX_RETRIES or not is_retryable (exc ):
71+ raise
72+ sleep_s = min (delay , _MAX_BACKOFF_S )
73+ _log .warning (
74+ "Transient gen-ai error (attempt %d/%d): %s; retrying in %.0fs" ,
75+ attempt + 1 ,
76+ _MAX_RETRIES + 1 ,
77+ exc ,
78+ sleep_s ,
79+ )
80+ time .sleep (sleep_s )
81+ delay *= _BACKOFF_FACTOR
82+ raise AssertionError ("unreachable" ) # loop either returns or raises
83+
2684
2785@dataclass
2886class _SseAccumulator :
@@ -114,12 +172,23 @@ def parse_sse_lines(lines: Iterable[str]) -> ChatResult:
114172 if not line or line .startswith ("event: " ) or not line .startswith (SSE_DATA_PREFIX ):
115173 continue
116174 data_str = line [len (SSE_DATA_PREFIX ) :]
175+ if _METADATA_SYNC_MARKER in data_str :
176+ raise TransientChatError (
177+ f"SSE transient error: { _METADATA_SYNC_MARKER } " ,
178+ status_code = None ,
179+ detail = None ,
180+ )
117181 try :
118182 event_data = json .loads (data_str )
119183 except json .JSONDecodeError :
120184 continue
121185 if "statusCode" in event_data :
122- raise RuntimeError (f"SSE error { event_data .get ('statusCode' )} : { event_data .get ('detail' )} " )
186+ code = event_data .get ("statusCode" )
187+ detail = event_data .get ("detail" )
188+ message = f"SSE error { code } : { detail } "
189+ if code in _RETRYABLE_STATUS_CODES :
190+ raise TransientChatError (message , status_code = code , detail = detail )
191+ raise ChatError (message , status_code = code , detail = detail )
123192 item = event_data .get ("item" )
124193 if not item :
125194 continue
@@ -149,12 +218,17 @@ def __init__(self, host: str, token: str, workspace_id: str, *, timeout: float =
149218 self ._client = httpx .Client (timeout = timeout )
150219
151220 def create_conversation (self ) -> str :
152- resp = self ._client .post (self ._base , headers = {** self ._auth , "Content-Type" : "application/json" })
153- resp .raise_for_status ()
154- body = resp .json ()
155- if "conversationId" not in body :
156- raise ValueError (f"GoodData /chat/conversations response missing 'conversationId': { body } " )
157- return body ["conversationId" ]
221+ def _do () -> str :
222+ resp = self ._client .post (self ._base , headers = {** self ._auth , "Content-Type" : "application/json" })
223+ resp .raise_for_status ()
224+ body = resp .json ()
225+ if "conversationId" not in body :
226+ raise ValueError (f"GoodData /chat/conversations response missing 'conversationId': { body } " )
227+ return body ["conversationId" ]
228+
229+ # NOTE: retrying create is not idempotent — a created-then-503 can leak an
230+ # orphaned (ephemeral) conversation. Acceptable for eval; do not reuse blindly.
231+ return _retry_transient (_do , is_retryable = _is_retryable_exc )
158232
159233 def delete_conversation (self , conversation_id : str ) -> None :
160234 try :
@@ -166,9 +240,13 @@ def send_message(self, conversation_id: str, question: str) -> ChatResult:
166240 url = f"{ self ._base } /{ conversation_id } /messages"
167241 headers = {** self ._auth , "Accept" : "text/event-stream" , "Content-Type" : "application/json" }
168242 body = {"item" : {"role" : "user" , "content" : {"type" : "text" , "text" : question }}}
169- with self ._client .stream ("POST" , url , json = body , headers = headers ) as resp :
170- resp .raise_for_status ()
171- return parse_sse_lines (resp .iter_lines ())
243+
244+ def _do () -> ChatResult :
245+ with self ._client .stream ("POST" , url , json = body , headers = headers ) as resp :
246+ resp .raise_for_status ()
247+ return parse_sse_lines (resp .iter_lines ())
248+
249+ return _retry_transient (_do , is_retryable = _is_retryable_exc )
172250
173251 def ask (self , item : DatasetItem ) -> ChatResult :
174252 """Run one single-turn conversation: create, send, parse, clean up."""
0 commit comments