Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
model_manager.py537 linesDownload Raw Back to core
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