Merge branch 'main' into docs-1.7-release

This commit is contained in:
Mendon Kissling
2025-12-02 13:30:16 -05:00
3 changed files with 226 additions and 123 deletions

View File

@ -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": {

View File

@ -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

View File

@ -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.")