MCPcopy Create free account
hub / github.com/Open-Bee/DataStudio / process_single_message

Method process_single_message

datastudio/models/base.py:343–396  ·  view source on GitHub ↗

Process a single message with retry logic. Args: message: Single processed message. all_kwargs: Arguments passed to generate_inner. Returns: Response content or fail_msg on failure.

(self, message, all_kwargs)

Source from the content-addressed store, hash-verified

341 return rd.random() * exp_delay
342
343 def process_single_message(self, message, all_kwargs):
344 """Process a single message with retry logic.
345
346 Args:
347 message: Single processed message.
348 all_kwargs: Arguments passed to generate_inner.
349
350 Returns:
351 Response content or fail_msg on failure.
352 """
353 try:
354 message = self.pre_process(message)
355 except Exception as e:
356 if self.logger:
357 self.logger.warning(f"Preprocessing failed: {str(e)}")
358 return self.fail_msg
359
360 # Initialize variables to avoid UnboundLocalError
361 ret_code = None
362 response_struct = None
363 response = None
364
365 for i in range(self.retry):
366 try:
367 ret_code, response_struct, response = self.generate_inner(
368 message, **all_kwargs
369 )
370
371 if ret_code == 0 and response_struct and response_struct != "":
372 return response_struct
373 else:
374 raise Exception("Invalid response")
375 except Exception as e:
376 # Use exponential backoff with jitter (faster than fixed wait)
377 if i < self.retry - 1:
378 delay = self._calculate_retry_delay(i)
379 time.sleep(delay)
380
381 if i + 1 > self.retry * 2 // 3 and self.logger:
382 self.logger.warning(
383 f"Attempt {i+1}/{self.retry} failed, RetCode: {ret_code}, "
384 f"Response: {response_struct}, Error: {type(e).__name__}: {str(e)[:200]}"
385 )
386
387 response_str = (
388 response.text
389 if hasattr(response, "text") and response
390 else str(response) if response else "No response"
391 )
392 if self.logger:
393 self.logger.warning(
394 f"Failed after {self.retry} attempts, RetCode: {ret_code}, Response: {response_str}"
395 )
396 return self.fail_msg
397
398 def pre_process(self, msg):
399 """Validate and preprocess raw input messages into API-ready format.

Callers 2

test_retry_on_failureMethod · 0.80

Calls 4

pre_processMethod · 0.95
generate_innerMethod · 0.95
warningMethod · 0.80

Tested by 2

test_retry_on_failureMethod · 0.64