Underground-Digital/Workflow-Engine
0
1import logging2import os3from collections.abc import Callable, Generator, Iterable, Sequence4from typing import IO, Any, Optional, Union, cast5 6from core.entities.embedding_type import EmbeddingInputType7from core.entities.provider_configuration import ProviderConfiguration, ProviderModelBundle8from core.entities.provider_entities import ModelLoadBalancingConfiguration9from core.errors.error import ProviderTokenNotInitError10from core.model_runtime.callbacks.base_callback import Callback11from core.model_runtime.entities.llm_entities import LLMResult12from core.model_runtime.entities.message_entities import PromptMessage, PromptMessageTool13from core.model_runtime.entities.model_entities import ModelType14from core.model_runtime.entities.rerank_entities import RerankResult15from core.model_runtime.entities.text_embedding_entities import TextEmbeddingResult16from core.model_runtime.errors.invoke import InvokeAuthorizationError, InvokeConnectionError, InvokeRateLimitError17from core.model_runtime.model_providers.__base.large_language_model import LargeLanguageModel18from core.model_runtime.model_providers.__base.moderation_model import ModerationModel19from core.model_runtime.model_providers.__base.rerank_model import RerankModel20from core.model_runtime.model_providers.__base.speech2text_model import Speech2TextModel21from core.model_runtime.model_providers.__base.text_embedding_model import TextEmbeddingModel22from core.model_runtime.model_providers.__base.tts_model import TTSModel23from core.provider_manager import ProviderManager24from extensions.ext_redis import redis_client25from models.provider import ProviderType26 27logger = logging.getLogger(__name__)28 29 30class ModelInstance:31 """32 Model instance class33 """34 35 def __init__(self, provider_model_bundle: ProviderModelBundle, model: str) -> None:36 self.provider_model_bundle = provider_model_bundle37 self.model = model38 self.provider = provider_model_bundle.configuration.provider.provider39 self.credentials = self._fetch_credentials_from_bundle(provider_model_bundle, model)40 self.model_type_instance = self.provider_model_bundle.model_type_instance41 self.load_balancing_manager = self._get_load_balancing_manager(42 configuration=provider_model_bundle.configuration,43 model_type=provider_model_bundle.model_type_instance.model_type,44 model=model,45 credentials=self.credentials,46 )47 48 @staticmethod49 def _fetch_credentials_from_bundle(provider_model_bundle: ProviderModelBundle, model: str) -> dict:50 """51 Fetch credentials from provider model bundle52 :param provider_model_bundle: provider model bundle53 :param model: model name54 :return:55 """56 configuration = provider_model_bundle.configuration57 model_type = provider_model_bundle.model_type_instance.model_type58 credentials = configuration.get_current_credentials(model_type=model_type, model=model)59 60 if credentials is None:61 raise ProviderTokenNotInitError(f"Model {model} credentials is not initialized.")62 63 return credentials64 65 @staticmethod66 def _get_load_balancing_manager(67 configuration: ProviderConfiguration, model_type: ModelType, model: str, credentials: dict68 ) -> Optional["LBModelManager"]:69 """70 Get load balancing model credentials71 :param configuration: provider configuration72 :param model_type: model type73 :param model: model name74 :param credentials: model credentials75 :return:76 """77 if configuration.model_settings and configuration.using_provider_type == ProviderType.CUSTOM:78 current_model_setting = None79 # check if model is disabled by admin80 for model_setting in configuration.model_settings:81 if model_setting.model_type == model_type and model_setting.model == model:82 current_model_setting = model_setting83 break84 85 # check if load balancing is enabled86 if current_model_setting and current_model_setting.load_balancing_configs:87 # use load balancing proxy to choose credentials88 lb_model_manager = LBModelManager(89 tenant_id=configuration.tenant_id,90 provider=configuration.provider.provider,91 model_type=model_type,92 model=model,93 load_balancing_configs=current_model_setting.load_balancing_configs,94 managed_credentials=credentials if configuration.custom_configuration.provider else None,95 )96 97 return lb_model_manager98 99 return None100 101 def invoke_llm(102 self,103 prompt_messages: list[PromptMessage],104 model_parameters: Optional[dict] = None,105 tools: Sequence[PromptMessageTool] | None = None,106 stop: Optional[list[str]] = None,107 stream: bool = True,108 user: Optional[str] = None,109 callbacks: Optional[list[Callback]] = None,110 ) -> Union[LLMResult, Generator]:111 """112 Invoke large language model113 114 :param prompt_messages: prompt messages115 :param model_parameters: model parameters116 :param tools: tools for tool calling117 :param stop: stop words118 :param stream: is stream response119 :param user: unique user id120 :param callbacks: callbacks121 :return: full response or stream response chunk generator result122 """123 if not isinstance(self.model_type_instance, LargeLanguageModel):124 raise Exception("Model type instance is not LargeLanguageModel")125 126 self.model_type_instance = cast(LargeLanguageModel, self.model_type_instance)127 return self._round_robin_invoke(128 function=self.model_type_instance.invoke,129 model=self.model,130 credentials=self.credentials,131 prompt_messages=prompt_messages,132 model_parameters=model_parameters,133 tools=tools,134 stop=stop,135 stream=stream,136 user=user,137 callbacks=callbacks,138 )139 140 def get_llm_num_tokens(141 self, prompt_messages: list[PromptMessage], tools: Optional[list[PromptMessageTool]] = None142 ) -> int:143 """144 Get number of tokens for llm145 146 :param prompt_messages: prompt messages147 :param tools: tools for tool calling148 :return:149 """150 if not isinstance(self.model_type_instance, LargeLanguageModel):151 raise Exception("Model type instance is not LargeLanguageModel")152 153 self.model_type_instance = cast(LargeLanguageModel, self.model_type_instance)154 return self._round_robin_invoke(155 function=self.model_type_instance.get_num_tokens,156 model=self.model,157 credentials=self.credentials,158 prompt_messages=prompt_messages,159 tools=tools,160 )161 162 def invoke_text_embedding(163 self, texts: list[str], user: Optional[str] = None, input_type: EmbeddingInputType = EmbeddingInputType.DOCUMENT164 ) -> TextEmbeddingResult:165 """166 Invoke large language model167 168 :param texts: texts to embed169 :param user: unique user id170 :param input_type: input type171 :return: embeddings result172 """173 if not isinstance(self.model_type_instance, TextEmbeddingModel):174 raise Exception("Model type instance is not TextEmbeddingModel")175 176 self.model_type_instance = cast(TextEmbeddingModel, self.model_type_instance)177 return self._round_robin_invoke(178 function=self.model_type_instance.invoke,179 model=self.model,180 credentials=self.credentials,181 texts=texts,182 user=user,183 input_type=input_type,184 )185 186 def get_text_embedding_num_tokens(self, texts: list[str]) -> int:187 """188 Get number of tokens for text embedding189 190 :param texts: texts to embed191 :return:192 """193 if not isinstance(self.model_type_instance, TextEmbeddingModel):194 raise Exception("Model type instance is not TextEmbeddingModel")195 196 self.model_type_instance = cast(TextEmbeddingModel, self.model_type_instance)197 return self._round_robin_invoke(198 function=self.model_type_instance.get_num_tokens,199 model=self.model,200 credentials=self.credentials,201 texts=texts,202 )203 204 def invoke_rerank(205 self,206 query: str,207 docs: list[str],208 score_threshold: Optional[float] = None,209 top_n: Optional[int] = None,210 user: Optional[str] = None,211 ) -> RerankResult:212 """213 Invoke rerank model214 215 :param query: search query216 :param docs: docs for reranking217 :param score_threshold: score threshold218 :param top_n: top n219 :param user: unique user id220 :return: rerank result221 """222 if not isinstance(self.model_type_instance, RerankModel):223 raise Exception("Model type instance is not RerankModel")224 225 self.model_type_instance = cast(RerankModel, self.model_type_instance)226 return self._round_robin_invoke(227 function=self.model_type_instance.invoke,228 model=self.model,229 credentials=self.credentials,230 query=query,231 docs=docs,232 score_threshold=score_threshold,233 top_n=top_n,234 user=user,235 )236 237 def invoke_moderation(self, text: str, user: Optional[str] = None) -> bool:238 """239 Invoke moderation model240 241 :param text: text to moderate242 :param user: unique user id243 :return: false if text is safe, true otherwise244 """245 if not isinstance(self.model_type_instance, ModerationModel):246 raise Exception("Model type instance is not ModerationModel")247 248 self.model_type_instance = cast(ModerationModel, self.model_type_instance)249 return self._round_robin_invoke(250 function=self.model_type_instance.invoke,251 model=self.model,252 credentials=self.credentials,253 text=text,254 user=user,255 )256 257 def invoke_speech2text(self, file: IO[bytes], user: Optional[str] = None) -> str:258 """259 Invoke large language model260 261 :param file: audio file262 :param user: unique user id263 :return: text for given audio file264 """265 if not isinstance(self.model_type_instance, Speech2TextModel):266 raise Exception("Model type instance is not Speech2TextModel")267 268 self.model_type_instance = cast(Speech2TextModel, self.model_type_instance)269 return self._round_robin_invoke(270 function=self.model_type_instance.invoke,271 model=self.model,272 credentials=self.credentials,273 file=file,274 user=user,275 )276 277 def invoke_tts(self, content_text: str, tenant_id: str, voice: str, user: Optional[str] = None) -> Iterable[bytes]:278 """279 Invoke large language tts model280 281 :param content_text: text content to be translated282 :param tenant_id: user tenant id283 :param voice: model timbre284 :param user: unique user id285 :return: text for given audio file286 """287 if not isinstance(self.model_type_instance, TTSModel):288 raise Exception("Model type instance is not TTSModel")289 290 self.model_type_instance = cast(TTSModel, self.model_type_instance)291 return self._round_robin_invoke(292 function=self.model_type_instance.invoke,293 model=self.model,294 credentials=self.credentials,295 content_text=content_text,296 user=user,297 tenant_id=tenant_id,298 voice=voice,299 )300 301 def _round_robin_invoke(self, function: Callable[..., Any], *args, **kwargs):302 """303 Round-robin invoke304 :param function: function to invoke305 :param args: function args306 :param kwargs: function kwargs307 :return:308 """309 if not self.load_balancing_manager:310 return function(*args, **kwargs)311 312 last_exception = None313 while True:314 lb_config = self.load_balancing_manager.fetch_next()315 if not lb_config:316 if not last_exception:317 raise ProviderTokenNotInitError("Model credentials is not initialized.")318 else:319 raise last_exception320 321 try:322 if "credentials" in kwargs:323 del kwargs["credentials"]324 return function(*args, **kwargs, credentials=lb_config.credentials)325 except InvokeRateLimitError as e:326 # expire in 60 seconds327 self.load_balancing_manager.cooldown(lb_config, expire=60)328 last_exception = e329 continue330 except (InvokeAuthorizationError, InvokeConnectionError) as e:331 # expire in 10 seconds332 self.load_balancing_manager.cooldown(lb_config, expire=10)333 last_exception = e334 continue335 except Exception as e:336 raise e337 338 def get_tts_voices(self, language: Optional[str] = None) -> list:339 """340 Invoke large language tts model voices341 342 :param language: tts language343 :return: tts model voices344 """345 if not isinstance(self.model_type_instance, TTSModel):346 raise Exception("Model type instance is not TTSModel")347 348 self.model_type_instance = cast(TTSModel, self.model_type_instance)349 return self.model_type_instance.get_tts_model_voices(350 model=self.model, credentials=self.credentials, language=language351 )352 353 354class ModelManager:355 def __init__(self) -> None:356 self._provider_manager = ProviderManager()357 358 def get_model_instance(self, tenant_id: str, provider: str, model_type: ModelType, model: str) -> ModelInstance:359 """360 Get model instance361 :param tenant_id: tenant id362 :param provider: provider name363 :param model_type: model type364 :param model: model name365 :return:366 """367 if not provider:368 return self.get_default_model_instance(tenant_id, model_type)369 370 provider_model_bundle = self._provider_manager.get_provider_model_bundle(371 tenant_id=tenant_id, provider=provider, model_type=model_type372 )373 374 return ModelInstance(provider_model_bundle, model)375 376 def get_default_provider_model_name(self, tenant_id: str, model_type: ModelType) -> tuple[str, str]:377 """378 Return first provider and the first model in the provider379 :param tenant_id: tenant id380 :param model_type: model type381 :return: provider name, model name382 """383 return self._provider_manager.get_first_provider_first_model(tenant_id, model_type)384 385 def get_default_model_instance(self, tenant_id: str, model_type: ModelType) -> ModelInstance:386 """387 Get default model instance388 :param tenant_id: tenant id389 :param model_type: model type390 :return:391 """392 default_model_entity = self._provider_manager.get_default_model(tenant_id=tenant_id, model_type=model_type)393 394 if not default_model_entity:395 raise ProviderTokenNotInitError(f"Default model not found for {model_type}")396 397 return self.get_model_instance(398 tenant_id=tenant_id,399 provider=default_model_entity.provider.provider,400 model_type=model_type,401 model=default_model_entity.model,402 )403 404 405class LBModelManager:406 def __init__(407 self,408 tenant_id: str,409 provider: str,410 model_type: ModelType,411 model: str,412 load_balancing_configs: list[ModelLoadBalancingConfiguration],413 managed_credentials: Optional[dict] = None,414 ) -> None:415 """416 Load balancing model manager417 :param tenant_id: tenant_id418 :param provider: provider419 :param model_type: model_type420 :param model: model name421 :param load_balancing_configs: all load balancing configurations422 :param managed_credentials: credentials if load balancing configuration name is __inherit__423 """424 self._tenant_id = tenant_id425 self._provider = provider426 self._model_type = model_type427 self._model = model428 self._load_balancing_configs = load_balancing_configs429 430 for load_balancing_config in self._load_balancing_configs[:]: # Iterate over a shallow copy of the list431 if load_balancing_config.name == "__inherit__":432 if not managed_credentials:433 # remove __inherit__ if managed credentials is not provided434 self._load_balancing_configs.remove(load_balancing_config)435 else:436 load_balancing_config.credentials = managed_credentials437 438 def fetch_next(self) -> Optional[ModelLoadBalancingConfiguration]:439 """440 Get next model load balancing config441 Strategy: Round Robin442 :return:443 """444 cache_key = "model_lb_index:{}:{}:{}:{}".format(445 self._tenant_id, self._provider, self._model_type.value, self._model446 )447 448 cooldown_load_balancing_configs = []449 max_index = len(self._load_balancing_configs)450 451 while True:452 current_index = redis_client.incr(cache_key)453 current_index = cast(int, current_index)454 if current_index >= 10000000:455 current_index = 1456 redis_client.set(cache_key, current_index)457 458 redis_client.expire(cache_key, 3600)459 if current_index > max_index:460 current_index = current_index % max_index461 462 real_index = current_index - 1463 if real_index > max_index:464 real_index = 0465 466 config = self._load_balancing_configs[real_index]467 468 if self.in_cooldown(config):469 cooldown_load_balancing_configs.append(config)470 if len(cooldown_load_balancing_configs) >= len(self._load_balancing_configs):471 # all configs are in cooldown472 return None473 474 continue475 476 if bool(os.environ.get("DEBUG", "False").lower() == "true"):477 logger.info(478 f"Model LB\nid: {config.id}\nname:{config.name}\n"479 f"tenant_id: {self._tenant_id}\nprovider: {self._provider}\n"480 f"model_type: {self._model_type.value}\nmodel: {self._model}"481 )482 483 return config484 485 return None486 487 def cooldown(self, config: ModelLoadBalancingConfiguration, expire: int = 60) -> None:488 """489 Cooldown model load balancing config490 :param config: model load balancing config491 :param expire: cooldown time492 :return:493 """494 cooldown_cache_key = "model_lb_index:cooldown:{}:{}:{}:{}:{}".format(495 self._tenant_id, self._provider, self._model_type.value, self._model, config.id496 )497 498 redis_client.setex(cooldown_cache_key, expire, "true")499 500 def in_cooldown(self, config: ModelLoadBalancingConfiguration) -> bool:501 """502 Check if model load balancing config is in cooldown503 :param config: model load balancing config504 :return:505 """506 cooldown_cache_key = "model_lb_index:cooldown:{}:{}:{}:{}:{}".format(507 self._tenant_id, self._provider, self._model_type.value, self._model, config.id508 )509 510 res = redis_client.exists(cooldown_cache_key)511 res = cast(bool, res)512 return res513 514 @staticmethod515 def get_config_in_cooldown_and_ttl(516 tenant_id: str, provider: str, model_type: ModelType, model: str, config_id: str517 ) -> tuple[bool, int]:518 """519 Get model load balancing config is in cooldown and ttl520 :param tenant_id: workspace id521 :param provider: provider name522 :param model_type: model type523 :param model: model name524 :param config_id: model load balancing config id525 :return:526 """527 cooldown_cache_key = "model_lb_index:cooldown:{}:{}:{}:{}:{}".format(528 tenant_id, provider, model_type.value, model, config_id529 )530 531 ttl = redis_client.ttl(cooldown_cache_key)532 if ttl == -2:533 return False, 0534 535 ttl = cast(int, ttl)536 return True, ttl537 