diff --git a/docs/openapi/openapi.json b/docs/openapi/openapi.json index f117ddde04..1fe418f702 100644 --- a/docs/openapi/openapi.json +++ b/docs/openapi/openapi.json @@ -2,7 +2,7 @@ "openapi": "3.1.0", "info": { "title": "Langflow", - "version": "1.6.8" + "version": "1.6.9" }, "paths": { "/api/v1/build/{flow_id}/vertices": { diff --git a/src/lfx/src/lfx/components/elastic/opensearch_multimodal.py b/src/lfx/src/lfx/components/elastic/opensearch_multimodal.py index 99a064d87b..29d0c515db 100644 --- a/src/lfx/src/lfx/components/elastic/opensearch_multimodal.py +++ b/src/lfx/src/lfx/components/elastic/opensearch_multimodal.py @@ -641,8 +641,14 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon @check_cached_vector_store def build_vector_store(self) -> OpenSearch: # Return raw OpenSearch client as our "vector store." - self.log(self.ingest_data) client = self.build_client() + + # Check if we're in ingestion-only mode (no search query) + has_search_query = bool((self.search_query or "").strip()) + if not has_search_query: + logger.debug("🔄 Ingestion-only mode activated: search operations will be skipped") + logger.debug("Starting ingestion mode...") + logger.warning(f"Embedding: {self.embedding}") self._add_documents_to_vector_store(client=client) return client @@ -660,25 +666,41 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon Args: client: OpenSearch client for performing operations """ + logger.debug("[INGESTION] _add_documents_to_vector_store called") # Convert DataFrame to Data if needed using parent's method self.ingest_data = self._prepare_ingest_data() + logger.debug( + f"[INGESTION] ingest_data type: " + f"{type(self.ingest_data)}, length: {len(self.ingest_data) if self.ingest_data else 0}" + ) + logger.debug( + f"[INGESTION] ingest_data content: " + f"{self.ingest_data[:2] if self.ingest_data and len(self.ingest_data) > 0 else 'empty'}" + ) + docs = self.ingest_data or [] if not docs: - self.log("No documents to ingest.") + logger.debug("✓ Ingestion complete: No documents provided") return if not self.embedding: msg = "Embedding handle is required to embed documents." raise ValueError(msg) - # Normalize embedding to list + # Normalize embedding to list first embeddings_list = self.embedding if isinstance(self.embedding, list) else [self.embedding] - if not embeddings_list: - msg = "At least one embedding is required to embed documents." - raise ValueError(msg) + # Filter out None values (fail-safe mode) - do this BEFORE checking if empty + embeddings_list = [e for e in embeddings_list if e is not None] + # NOW check if we have any valid embeddings left after filtering + if not embeddings_list: + logger.warning("All embeddings returned None (fail-safe mode enabled). Skipping document ingestion.") + self.log("Embedding returned None (fail-safe mode enabled). Skipping document ingestion.") + return + + logger.debug(f"[INGESTION] Valid embeddings after filtering: {len(embeddings_list)}") self.log(f"Available embedding models: {len(embeddings_list)}") # Select the embedding to use for ingestion @@ -790,6 +812,7 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon dynamic_field_name = get_embedding_field_name(embedding_model) + logger.info(f"✓ Selected embedding model for ingestion: '{embedding_model}'") self.log(f"Using embedding model for ingestion: {embedding_model}") self.log(f"Dynamic vector field: {dynamic_field_name}") @@ -814,6 +837,7 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon metadatas = [] # Process docs_metadata table input into a dict additional_metadata = {} + logger.debug(f"[LF] Docs metadata {self.docs_metadata}") if hasattr(self, "docs_metadata") and self.docs_metadata: logger.info(f"[LF] Docs metadata {self.docs_metadata}") if isinstance(self.docs_metadata[-1], Data): @@ -956,6 +980,9 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon ) self.log(metadatas) + logger.info( + f"✓ Ingestion complete: Successfully indexed {len(return_ids)} documents with model '{embedding_model}'" + ) self.log(f"Successfully indexed {len(return_ids)} documents with model {embedding_model}.") # ---------- helpers for filters ---------- @@ -1172,6 +1199,11 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon msg = "Embedding is required to run hybrid search (KNN + keyword)." raise ValueError(msg) + # Check if embedding is None (fail-safe mode) + if self.embedding is None or (isinstance(self.embedding, list) and all(e is None for e in self.embedding)): + logger.error("Embedding returned None (fail-safe mode enabled). Cannot perform search.") + return [] + # Build filter clauses first so we can use them in model detection filter_clauses = self._coerce_filter_clauses(filter_obj) @@ -1187,6 +1219,14 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon # Normalize embedding to list embeddings_list = self.embedding if isinstance(self.embedding, list) else [self.embedding] + # Filter out None values (fail-safe mode) + embeddings_list = [e for e in embeddings_list if e is not None] + + if not embeddings_list: + logger.error( + "No valid embeddings available after filtering None values (fail-safe mode). Cannot perform search." + ) + return [] # Create a comprehensive map of model names to embedding objects # Check all possible identifiers (deployment, model, model_id, model_name) @@ -1518,6 +1558,9 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon This is the main interface method that performs the multi-model search using the configured search_query and returns results in Langflow's Data format. + Always builds the vector store (triggering ingestion if needed), then performs + search only if a query is provided. + Returns: List of Data objects containing search results with text and metadata @@ -1525,9 +1568,19 @@ class OpenSearchVectorStoreComponentMultimodalMultiEmbedding(LCVectorStoreCompon Exception: If search operation fails """ try: - raw = self.search(self.search_query or "") + # Always build/cache the vector store to ensure ingestion happens + if self._cached_vector_store is None: + self.build_vector_store() + + # Only perform search if query is provided + search_query = (self.search_query or "").strip() + if not search_query: + self.log("No search query provided - ingestion completed, returning empty results") + return [] + + # Perform search with the provided query + raw = self.search(search_query) return [Data(text=hit["page_content"], **hit["metadata"]) for hit in raw] - self.log(self.ingest_data) except Exception as e: self.log(f"search_documents error: {e}") raise diff --git a/src/lfx/src/lfx/components/models_and_agents/embedding_model.py b/src/lfx/src/lfx/components/models_and_agents/embedding_model.py index 62838687b3..1b436ac74f 100644 --- a/src/lfx/src/lfx/components/models_and_agents/embedding_model.py +++ b/src/lfx/src/lfx/components/models_and_agents/embedding_model.py @@ -132,6 +132,15 @@ class EmbeddingModelComponent(LCEmbeddingsModel): advanced=True, show=False, ), + BoolInput( + name="fail_safe_mode", + display_name="Fail-Safe Mode", + value=False, + advanced=True, + info="When enabled, errors will be logged instead of raising exceptions. " + "The component will return None on error.", + real_time_refresh=True, + ), ] @staticmethod @@ -152,6 +161,19 @@ class EmbeddingModelComponent(LCEmbeddingsModel): logger.exception("Error fetching models") return WATSONX_EMBEDDING_MODEL_NAMES + async def fetch_ollama_models(self) -> list[str]: + try: + return await get_ollama_models( + base_url_value=self.ollama_base_url, + desired_capability=DESIRED_CAPABILITY, + json_models_key=JSON_MODELS_KEY, + json_name_key=JSON_NAME_KEY, + json_capabilities_key=JSON_CAPABILITIES_KEY, + ) + except Exception: # noqa: BLE001 + logger.exception("Error fetching models") + return [] + async def build_embeddings(self) -> Embeddings: provider = self.provider model = self.model @@ -169,27 +191,16 @@ class EmbeddingModelComponent(LCEmbeddingsModel): if provider == "OpenAI": if not api_key: msg = "OpenAI API key is required when using OpenAI provider" + if self.fail_safe_mode: + logger.error(msg) + return None raise ValueError(msg) - # Create the primary embedding instance - embeddings_instance = OpenAIEmbeddings( - model=model, - dimensions=dimensions or None, - base_url=api_base or None, - api_key=api_key, - chunk_size=chunk_size, - max_retries=max_retries, - timeout=request_timeout or None, - show_progress_bar=show_progress_bar, - model_kwargs=model_kwargs, - ) - - # Create dedicated instances for each available model - available_models_dict = {} - for model_name in OPENAI_EMBEDDING_MODEL_NAMES: - available_models_dict[model_name] = OpenAIEmbeddings( - model=model_name, - dimensions=dimensions or None, # Use same dimensions config for all + try: + # Create the primary embedding instance + embeddings_instance = OpenAIEmbeddings( + model=model, + dimensions=dimensions or None, base_url=api_base or None, api_key=api_key, chunk_size=chunk_size, @@ -199,10 +210,31 @@ class EmbeddingModelComponent(LCEmbeddingsModel): model_kwargs=model_kwargs, ) - return EmbeddingsWithModels( - embeddings=embeddings_instance, - available_models=available_models_dict, - ) + # Create dedicated instances for each available model + available_models_dict = {} + for model_name in OPENAI_EMBEDDING_MODEL_NAMES: + available_models_dict[model_name] = OpenAIEmbeddings( + model=model_name, + dimensions=dimensions or None, # Use same dimensions config for all + base_url=api_base or None, + api_key=api_key, + chunk_size=chunk_size, + max_retries=max_retries, + timeout=request_timeout or None, + show_progress_bar=show_progress_bar, + model_kwargs=model_kwargs, + ) + + return EmbeddingsWithModels( + embeddings=embeddings_instance, + available_models=available_models_dict, + ) + except Exception as e: + msg = f"Failed to initialize OpenAI embeddings: {e}" + if self.fail_safe_mode: + logger.error(msg) + return None + raise if provider == "Ollama": try: @@ -212,124 +244,159 @@ class EmbeddingModelComponent(LCEmbeddingsModel): from langchain_community.embeddings import OllamaEmbeddings except ImportError: msg = "Please install langchain-ollama: pip install langchain-ollama" + if self.fail_safe_mode: + logger.error(msg) + return None raise ImportError(msg) from None - transformed_base_url = transform_localhost_url(ollama_base_url) + try: + transformed_base_url = transform_localhost_url(ollama_base_url) - # Check if URL contains /v1 suffix (OpenAI-compatible mode) - if transformed_base_url and transformed_base_url.rstrip("/").endswith("/v1"): - # Strip /v1 suffix and log warning - transformed_base_url = transformed_base_url.rstrip("/").removesuffix("/v1") - logger.warning( - "Detected '/v1' suffix in base URL. The Ollama component uses the native Ollama API, " - "not the OpenAI-compatible API. The '/v1' suffix has been automatically removed. " - "If you want to use the OpenAI-compatible API, please use the OpenAI component instead. " - "Learn more at https://docs.ollama.com/openai#openai-compatibility" - ) + # Check if URL contains /v1 suffix (OpenAI-compatible mode) + if transformed_base_url and transformed_base_url.rstrip("/").endswith("/v1"): + # Strip /v1 suffix and log warning + transformed_base_url = transformed_base_url.rstrip("/").removesuffix("/v1") + logger.warning( + "Detected '/v1' suffix in base URL. The Ollama component uses the native Ollama API, " + "not the OpenAI-compatible API. The '/v1' suffix has been automatically removed. " + "If you want to use the OpenAI-compatible API, please use the OpenAI component instead. " + "Learn more at https://docs.ollama.com/openai#openai-compatibility" + ) - final_base_url = transformed_base_url or "http://localhost:11434" + final_base_url = transformed_base_url or "http://localhost:11434" - # Create the primary embedding instance - embeddings_instance = OllamaEmbeddings( - model=model, - base_url=final_base_url, - **model_kwargs, - ) - - # Fetch available Ollama models - available_model_names = await get_ollama_models( - base_url_value=self.ollama_base_url, - desired_capability=DESIRED_CAPABILITY, - json_models_key=JSON_MODELS_KEY, - json_name_key=JSON_NAME_KEY, - json_capabilities_key=JSON_CAPABILITIES_KEY, - ) - - # Create dedicated instances for each available model - available_models_dict = {} - for model_name in available_model_names: - available_models_dict[model_name] = OllamaEmbeddings( - model=model_name, + # Create the primary embedding instance + embeddings_instance = OllamaEmbeddings( + model=model, base_url=final_base_url, **model_kwargs, ) - return EmbeddingsWithModels( - embeddings=embeddings_instance, - available_models=available_models_dict, - ) + # Fetch available Ollama models + available_model_names = await self.fetch_ollama_models() + + # Create dedicated instances for each available model + available_models_dict = {} + for model_name in available_model_names: + available_models_dict[model_name] = OllamaEmbeddings( + model=model_name, + base_url=final_base_url, + **model_kwargs, + ) + + return EmbeddingsWithModels( + embeddings=embeddings_instance, + available_models=available_models_dict, + ) + except Exception as e: + msg = f"Failed to initialize Ollama embeddings: {e}" + if self.fail_safe_mode: + logger.error(msg) + return None + raise if provider == "IBM watsonx.ai": try: from langchain_ibm import WatsonxEmbeddings except ImportError: msg = "Please install langchain-ibm: pip install langchain-ibm" + if self.fail_safe_mode: + logger.error(msg) + return None raise ImportError(msg) from None if not api_key: msg = "IBM watsonx.ai API key is required when using IBM watsonx.ai provider" + if self.fail_safe_mode: + logger.error(msg) + return None raise ValueError(msg) project_id = self.project_id if not project_id: msg = "Project ID is required for IBM watsonx.ai provider" + if self.fail_safe_mode: + logger.error(msg) + return None raise ValueError(msg) - from ibm_watsonx_ai import APIClient, Credentials + try: + from ibm_watsonx_ai import APIClient, Credentials - final_url = base_url_ibm_watsonx or "https://us-south.ml.cloud.ibm.com" + final_url = base_url_ibm_watsonx or "https://us-south.ml.cloud.ibm.com" - credentials = Credentials( - api_key=self.api_key, - url=final_url, - ) + credentials = Credentials( + api_key=self.api_key, + url=final_url, + ) - api_client = APIClient(credentials) + api_client = APIClient(credentials) - params = { - EmbedTextParamsMetaNames.TRUNCATE_INPUT_TOKENS: self.truncate_input_tokens, - EmbedTextParamsMetaNames.RETURN_OPTIONS: {"input_text": self.input_text}, - } + params = { + EmbedTextParamsMetaNames.TRUNCATE_INPUT_TOKENS: self.truncate_input_tokens, + EmbedTextParamsMetaNames.RETURN_OPTIONS: {"input_text": self.input_text}, + } - # Create the primary embedding instance - embeddings_instance = WatsonxEmbeddings( - model_id=model, - params=params, - watsonx_client=api_client, - project_id=project_id, - ) - - # Fetch available IBM watsonx.ai models - available_model_names = self.fetch_ibm_models(final_url) - - # Create dedicated instances for each available model - available_models_dict = {} - for model_name in available_model_names: - available_models_dict[model_name] = WatsonxEmbeddings( - model_id=model_name, + # Create the primary embedding instance + embeddings_instance = WatsonxEmbeddings( + model_id=model, params=params, watsonx_client=api_client, project_id=project_id, ) - return EmbeddingsWithModels( - embeddings=embeddings_instance, - available_models=available_models_dict, - ) + # Fetch available IBM watsonx.ai models + available_model_names = self.fetch_ibm_models(final_url) + + # Create dedicated instances for each available model + available_models_dict = {} + for model_name in available_model_names: + available_models_dict[model_name] = WatsonxEmbeddings( + model_id=model_name, + params=params, + watsonx_client=api_client, + project_id=project_id, + ) + + return EmbeddingsWithModels( + embeddings=embeddings_instance, + available_models=available_models_dict, + ) + except Exception as e: + msg = f"Failed to authenticate with IBM watsonx.ai: {e}" + if self.fail_safe_mode: + logger.error(msg) + return None + raise msg = f"Unknown provider: {provider}" + if self.fail_safe_mode: + logger.error(msg) + return None raise ValueError(msg) async def update_build_config( self, build_config: dotdict, field_value: Any, field_name: str | None = None ) -> dotdict: + # Handle fail_safe_mode changes first - set all required fields to False if enabled + if field_name == "fail_safe_mode": + if field_value: # If fail_safe_mode is enabled + build_config["api_key"]["required"] = False + elif hasattr(self, "provider"): + # If fail_safe_mode is disabled, restore required flags based on provider + if self.provider in ["OpenAI", "IBM watsonx.ai"]: + build_config["api_key"]["required"] = True + else: # Ollama + build_config["api_key"]["required"] = False + if field_name == "provider": if field_value == "OpenAI": build_config["model"]["options"] = OPENAI_EMBEDDING_MODEL_NAMES build_config["model"]["value"] = OPENAI_EMBEDDING_MODEL_NAMES[0] build_config["api_key"]["display_name"] = "OpenAI API Key" - build_config["api_key"]["required"] = True + # Only set required=True if fail_safe_mode is not enabled + build_config["api_key"]["required"] = not (hasattr(self, "fail_safe_mode") and self.fail_safe_mode) build_config["api_key"]["show"] = True build_config["api_base"]["display_name"] = "OpenAI API Base URL" build_config["api_base"]["advanced"] = True @@ -344,13 +411,7 @@ class EmbeddingModelComponent(LCEmbeddingsModel): if await is_valid_ollama_url(url=self.ollama_base_url): try: - models = await get_ollama_models( - base_url_value=self.ollama_base_url, - desired_capability=DESIRED_CAPABILITY, - json_models_key=JSON_MODELS_KEY, - json_name_key=JSON_NAME_KEY, - json_capabilities_key=JSON_CAPABILITIES_KEY, - ) + models = await self.fetch_ollama_models() build_config["model"]["options"] = models build_config["model"]["value"] = models[0] if models else "" except ValueError: @@ -372,7 +433,8 @@ class EmbeddingModelComponent(LCEmbeddingsModel): build_config["model"]["options"] = self.fetch_ibm_models(base_url=self.base_url_ibm_watsonx) build_config["model"]["value"] = self.fetch_ibm_models(base_url=self.base_url_ibm_watsonx)[0] build_config["api_key"]["display_name"] = "IBM watsonx.ai API Key" - build_config["api_key"]["required"] = True + # Only set required=True if fail_safe_mode is not enabled + build_config["api_key"]["required"] = not (hasattr(self, "fail_safe_mode") and self.fail_safe_mode) build_config["api_key"]["show"] = True build_config["api_base"]["show"] = False build_config["ollama_base_url"]["show"] = False @@ -390,13 +452,7 @@ class EmbeddingModelComponent(LCEmbeddingsModel): ollama_url = self.ollama_base_url if await is_valid_ollama_url(url=ollama_url): try: - models = await get_ollama_models( - base_url_value=ollama_url, - desired_capability=DESIRED_CAPABILITY, - json_models_key=JSON_MODELS_KEY, - json_name_key=JSON_NAME_KEY, - json_capabilities_key=JSON_CAPABILITIES_KEY, - ) + models = await self.fetch_ollama_models() build_config["model"]["options"] = models build_config["model"]["value"] = models[0] if models else "" except ValueError: @@ -408,13 +464,7 @@ class EmbeddingModelComponent(LCEmbeddingsModel): ollama_url = self.ollama_base_url if await is_valid_ollama_url(url=ollama_url): try: - models = await get_ollama_models( - base_url_value=ollama_url, - desired_capability=DESIRED_CAPABILITY, - json_models_key=JSON_MODELS_KEY, - json_name_key=JSON_NAME_KEY, - json_capabilities_key=JSON_CAPABILITIES_KEY, - ) + models = await self.fetch_ollama_models() build_config["model"]["options"] = models except ValueError: await logger.awarning("Failed to refresh Ollama embedding models.")