diff --git a/.secrets.baseline b/.secrets.baseline index 4e12418279..235ad435ed 100644 --- a/.secrets.baseline +++ b/.secrets.baseline @@ -3365,15 +3365,6 @@ "is_secret": false } ], - "src/backend/tests/unit/components/files_and_knowledge/test_ingestion.py": [ - { - "type": "Secret Keyword", - "filename": "src/backend/tests/unit/components/files_and_knowledge/test_ingestion.py", - "hashed_secret": "3acfb2c2b433c0ea7ff107e33df91b18e52f960f", - "is_verified": false, - "line_number": 312 - } - ], "src/backend/tests/unit/components/files_and_knowledge/test_retrieval.py": [ { "type": "Secret Keyword", @@ -8414,5 +8405,5 @@ } ] }, - "generated_at": "2026-04-09T14:26:21Z" + "generated_at": "2026-04-10T17:09:16Z" } diff --git a/src/backend/tests/unit/components/files_and_knowledge/test_ingestion.py b/src/backend/tests/unit/components/files_and_knowledge/test_ingestion.py index a8e4e44fdc..461bf0b77e 100644 --- a/src/backend/tests/unit/components/files_and_knowledge/test_ingestion.py +++ b/src/backend/tests/unit/components/files_and_knowledge/test_ingestion.py @@ -106,6 +106,16 @@ class TestKnowledgeIngestionComponent(ComponentTestBaseWithClient): with pytest.raises(ValueError, match="Column 'nonexistent' not found in DataFrame"): component._validate_column_config(data_df) + def test_new_knowledge_dialog_uses_provider_credentials(self, component_class, default_kwargs): + """Test the create-knowledge dialog no longer exposes a redundant API key override.""" + component = component_class(**default_kwargs) + dialog_inputs = component.inputs[0].dialog_inputs["fields"]["data"]["node"] + embedding_model_input = dialog_inputs["template"]["02_embedding_model"] + + assert dialog_inputs["field_order"] == ["01_new_kb_name", "02_embedding_model"] + assert "03_api_key" not in dialog_inputs["template"] + assert "configured credentials" in embedding_model_input.info + @patch("lfx.components.files_and_knowledge.ingestion.get_settings_service") @patch("lfx.components.files_and_knowledge.ingestion.encrypt_api_key") def test_build_embedding_metadata(self, mock_encrypt, mock_get_settings, component_class, default_kwargs): @@ -309,7 +319,6 @@ class TestKnowledgeIngestionComponent(ComponentTestBaseWithClient): field_value = { "01_new_kb_name": "new_test_kb", "02_embedding_model": model_selection, - "03_api_key": "test-key", } # Mock embedding validation @@ -322,8 +331,8 @@ class TestKnowledgeIngestionComponent(ComponentTestBaseWithClient): assert result["knowledge_base"]["value"] == "new_test_kb" assert "new_test_kb" in result["knowledge_base"]["options"] - assert mock_get_embeddings.call_args.kwargs["api_key"] == "test-key" - assert mock_save_metadata.call_args.kwargs["api_key"] == "test-key" + assert "api_key" not in mock_get_embeddings.call_args.kwargs + assert "api_key" not in mock_save_metadata.call_args.kwargs @patch("lfx.components.files_and_knowledge.ingestion.get_embeddings") async def test_build_kb_info_with_message_input(self, mock_get_embeddings, component_class, default_kwargs): @@ -353,7 +362,6 @@ class TestKnowledgeIngestionComponent(ComponentTestBaseWithClient): field_value = { "01_new_kb_name": "invalid@name", # Invalid character "02_embedding_model": "sentence-transformers/all-MiniLM-L6-v2", - "03_api_key": None, } with pytest.raises(ValueError, match="Invalid knowledge base name"): diff --git a/src/lfx/src/lfx/_assets/component_index.json b/src/lfx/src/lfx/_assets/component_index.json index 8fac6cd7a5..5076f5698a 100644 --- a/src/lfx/src/lfx/_assets/component_index.json +++ b/src/lfx/src/lfx/_assets/component_index.json @@ -69138,7 +69138,7 @@ "icon": "upload", "legacy": false, "metadata": { - "code_hash": "4573b85351ee", + "code_hash": "1b571470646c", "dependencies": { "dependencies": [ { @@ -69270,7 +69270,7 @@ "show": true, "title_case": false, "type": "code", - "value": "from __future__ import annotations\n\nimport asyncio\nimport hashlib\nimport json\nimport re\nimport uuid\nfrom dataclasses import asdict, dataclass, field\nfrom datetime import datetime, timezone\nfrom pathlib import Path\nfrom typing import TYPE_CHECKING, Any\n\nimport pandas as pd\nfrom cryptography.fernet import InvalidToken\nfrom langchain_chroma import Chroma\nfrom langflow.services.auth.utils import decrypt_api_key, encrypt_api_key\nfrom langflow.services.database.models.user.crud import get_user_by_id\n\nfrom lfx.base.knowledge_bases.knowledge_base_utils import get_knowledge_bases\nfrom lfx.base.models.unified_models import get_embedding_model_options, get_embeddings\nfrom lfx.components.processing.converter import convert_to_dataframe\nfrom lfx.custom import Component\nfrom lfx.io import (\n BoolInput,\n DropdownInput,\n HandleInput,\n IntInput,\n ModelInput,\n Output,\n SecretStrInput,\n StrInput,\n TableInput,\n)\nfrom lfx.schema.data import Data\nfrom lfx.schema.table import EditMode\nfrom lfx.services.deps import (\n get_settings_service,\n session_scope,\n)\nfrom lfx.utils.validate_cloud import raise_error_if_astra_cloud_disable_component\n\nif TYPE_CHECKING:\n from lfx.schema.dataframe import DataFrame\n\n_KNOWLEDGE_BASES_ROOT_PATH: Path | None = None\n\n# Error message to raise if we're in Astra cloud environment and the component is not supported.\nastra_error_msg = \"Knowledge ingestion is not supported in Astra cloud environment.\"\n\n\ndef _get_knowledge_bases_root_path() -> Path:\n \"\"\"Lazy load the knowledge bases root path from settings.\"\"\"\n global _KNOWLEDGE_BASES_ROOT_PATH # noqa: PLW0603\n if _KNOWLEDGE_BASES_ROOT_PATH is None:\n settings = get_settings_service().settings\n knowledge_directory = settings.knowledge_bases_dir\n if not knowledge_directory:\n msg = \"Knowledge bases directory is not set in the settings.\"\n raise ValueError(msg)\n _KNOWLEDGE_BASES_ROOT_PATH = Path(knowledge_directory).expanduser()\n return _KNOWLEDGE_BASES_ROOT_PATH\n\n\nclass KnowledgeIngestionComponent(Component):\n \"\"\"Create or append to Langflow Knowledge from a DataFrame.\"\"\"\n\n # ------ UI metadata ---------------------------------------------------\n display_name = \"Knowledge Ingestion\"\n description = \"Create or update knowledge in Langflow.\"\n icon = \"upload\"\n name = \"KnowledgeIngestion\"\n\n def __init__(self, *args, **kwargs) -> None:\n super().__init__(*args, **kwargs)\n self._cached_kb_path: Path | None = None\n\n @dataclass\n class NewKnowledgeBaseInput:\n functionality: str = \"create\"\n fields: dict[str, dict] = field(\n default_factory=lambda: {\n \"data\": {\n \"node\": {\n \"name\": \"create_knowledge_base\",\n \"description\": \"Create new knowledge in Langflow.\",\n \"display_name\": \"Create new knowledge\",\n \"field_order\": [\n \"01_new_kb_name\",\n \"02_embedding_model\",\n \"03_api_key\",\n ],\n \"template\": {\n \"01_new_kb_name\": StrInput(\n name=\"new_kb_name\",\n display_name=\"Knowledge Name\",\n info=\"Name of the new knowledge to create.\",\n required=True,\n ),\n \"02_embedding_model\": ModelInput(\n name=\"embedding_model\",\n display_name=\"Choose Embedding Model\",\n info=\"Select the embedding model to use for this knowledge base.\",\n required=True,\n model_type=\"embedding\",\n ),\n \"03_api_key\": SecretStrInput(\n name=\"api_key\",\n display_name=\"Embedding Provider API Key\",\n info=\"Optional API key override used to validate and save this knowledge base.\",\n required=False,\n advanced=True,\n ),\n },\n },\n }\n }\n )\n\n # ------ Inputs --------------------------------------------------------\n inputs = [\n DropdownInput(\n name=\"knowledge_base\",\n display_name=\"Knowledge\",\n info=\"Select the knowledge to load data from.\",\n required=True,\n options=[],\n refresh_button=True,\n real_time_refresh=True,\n dialog_inputs=asdict(NewKnowledgeBaseInput()),\n ),\n HandleInput(\n name=\"input_df\",\n display_name=\"Input\",\n info=(\n \"Table with all original columns (already chunked / processed). \"\n \"Accepts Message, Data, or DataFrame. If Message or Data is provided, \"\n \"it is converted to a DataFrame automatically.\"\n ),\n input_types=[\"Message\", \"Data\", \"JSON\", \"DataFrame\", \"Table\"],\n required=True,\n ),\n TableInput(\n name=\"column_config\",\n display_name=\"Column Configuration\",\n info=\"Configure column behavior for the knowledge base.\",\n required=True,\n table_schema=[\n {\n \"name\": \"column_name\",\n \"display_name\": \"Column Name\",\n \"type\": \"str\",\n \"description\": \"Name of the column in the source DataFrame\",\n \"edit_mode\": EditMode.INLINE,\n },\n {\n \"name\": \"vectorize\",\n \"display_name\": \"Vectorize\",\n \"type\": \"boolean\",\n \"description\": \"Create embeddings for this column\",\n \"default\": False,\n \"edit_mode\": EditMode.INLINE,\n },\n {\n \"name\": \"identifier\",\n \"display_name\": \"Identifier\",\n \"type\": \"boolean\",\n \"description\": \"Use this column as unique identifier\",\n \"default\": False,\n \"edit_mode\": EditMode.INLINE,\n },\n ],\n value=[\n {\n \"column_name\": \"text\",\n \"vectorize\": True,\n \"identifier\": True,\n },\n ],\n ),\n IntInput(\n name=\"chunk_size\",\n display_name=\"Chunk Size\",\n info=\"Batch size for processing embeddings\",\n advanced=True,\n value=1000,\n ),\n SecretStrInput(\n name=\"api_key\",\n display_name=\"Embedding Provider API Key\",\n info=\"Overrides global provider settings. Leave blank to use your pre-configured API Key.\",\n advanced=True,\n required=False,\n ),\n BoolInput(\n name=\"allow_duplicates\",\n display_name=\"Allow Duplicates\",\n info=\"Allow duplicate rows in the knowledge base\",\n advanced=True,\n value=False,\n ),\n ]\n\n # ------ Outputs -------------------------------------------------------\n outputs = [Output(display_name=\"Results\", name=\"dataframe_output\", method=\"build_kb_info\")]\n\n # ------ Internal helpers ---------------------------------------------\n def _get_kb_root(self) -> Path:\n \"\"\"Return the root directory for knowledge bases.\"\"\"\n return _get_knowledge_bases_root_path()\n\n @staticmethod\n def _scalar_notna(value) -> bool:\n \"\"\"Check if a value is not NA, safely handling arrays and sequences.\n\n ``pd.notna`` returns an array when given an array-like input, which\n cannot be used directly in a boolean context. This helper collapses the\n result to a single scalar ``bool``.\n \"\"\"\n result = pd.notna(value)\n # If result is array-like (numpy array, list, etc.), treat non-empty arrays as \"present\"\n if hasattr(result, \"__iter__\") and not isinstance(result, str):\n import numpy as np\n\n arr = np.asarray(result)\n return arr.size > 0 and arr.all()\n return bool(result)\n\n def _validate_column_config(self, df_source: pd.DataFrame) -> list[dict[str, Any]]:\n \"\"\"Validate column configuration using Structured Output patterns.\"\"\"\n if not self.column_config:\n msg = \"Column configuration cannot be empty\"\n raise ValueError(msg)\n\n # Convert table input to list of dicts (similar to Structured Output)\n config_list = self.column_config if isinstance(self.column_config, list) else []\n\n # Validate column names exist in DataFrame\n df_columns = set(df_source.columns)\n for config in config_list:\n col_name = config.get(\"column_name\")\n if col_name not in df_columns:\n msg = f\"Column '{col_name}' not found in DataFrame. Available columns: {sorted(df_columns)}\"\n raise ValueError(msg)\n\n return config_list\n\n def _build_embedding_metadata(\n self,\n model_selection: list[dict[str, Any]],\n api_key: str | None = None,\n ) -> dict[str, Any]:\n \"\"\"Build embedding model metadata from a model selection dict.\n\n Args:\n model_selection: Model selection list from ModelInput\n (e.g. [{'name': ..., 'provider': ..., 'metadata': ...}])\n api_key: Optional API key override.\n \"\"\"\n model_dict = model_selection[0] if isinstance(model_selection, list) else model_selection\n embedding_model = model_dict.get(\"name\", \"\")\n embedding_provider = model_dict.get(\"provider\", \"Unknown\")\n\n api_key_to_save = None\n if api_key and hasattr(api_key, \"get_secret_value\"):\n api_key_to_save = api_key.get_secret_value()\n elif isinstance(api_key, str):\n api_key_to_save = api_key\n\n encrypted_api_key = None\n if api_key_to_save:\n settings_service = get_settings_service()\n try:\n encrypted_api_key = encrypt_api_key(api_key_to_save, settings_service=settings_service)\n except (TypeError, ValueError) as e:\n self.log(f\"Could not encrypt API key: {e}\")\n\n return {\n \"embedding_provider\": embedding_provider,\n \"embedding_model\": embedding_model,\n \"model_selection\": model_dict, # Store full selection for get_embeddings() reconstruction\n \"api_key\": encrypted_api_key,\n \"api_key_used\": bool(api_key),\n \"chunk_size\": self.chunk_size,\n \"created_at\": datetime.now(timezone.utc).isoformat(),\n }\n\n def _save_embedding_metadata(\n self,\n kb_path: Path,\n model_selection: list[dict[str, Any]],\n api_key: str | None = None,\n ) -> None:\n \"\"\"Save embedding model metadata.\"\"\"\n embedding_metadata = self._build_embedding_metadata(model_selection, api_key)\n metadata_path = kb_path / \"embedding_metadata.json\"\n metadata_path.write_text(json.dumps(embedding_metadata, indent=2))\n\n def _update_metadata_metrics(self, kb_path: Path, chroma: Chroma) -> None:\n \"\"\"Update embedding_metadata.json with accurate chunk/word/character counts.\n\n This ensures the Knowledge Base modal displays correct stats after\n component-based ingestion, matching the behavior of API-based ingestion.\n Delegates to KBAnalysisHelper.update_text_metrics to avoid duplicating\n the batched metrics counting logic.\n \"\"\"\n import chromadb.errors\n from langflow.api.utils.kb_helpers import KBAnalysisHelper, KBStorageHelper\n\n metadata_path = kb_path / \"embedding_metadata.json\"\n if not metadata_path.exists():\n return\n\n try:\n metadata = json.loads(metadata_path.read_text())\n KBAnalysisHelper.update_text_metrics(kb_path, metadata, chroma)\n metadata[\"size\"] = KBStorageHelper.get_directory_size(kb_path)\n metadata_path.write_text(json.dumps(metadata, indent=2))\n except (OSError, ValueError, TypeError, json.JSONDecodeError, chromadb.errors.ChromaError) as e:\n self.log(f\"Warning: Could not update metadata metrics: {e}\")\n\n def _save_kb_files(\n self,\n kb_path: Path,\n config_list: list[dict[str, Any]],\n ) -> None:\n \"\"\"Save KB files using File Component storage patterns.\"\"\"\n try:\n # Create directory (following File Component patterns)\n kb_path.mkdir(parents=True, exist_ok=True)\n\n # Save column configuration\n # Only do this if the file doesn't exist already\n cfg_path = kb_path / \"schema.json\"\n if not cfg_path.exists():\n cfg_path.write_text(json.dumps(config_list, indent=2))\n\n except (OSError, TypeError, ValueError) as e:\n self.log(f\"Error saving KB files: {e}\")\n\n def _build_column_metadata(self, config_list: list[dict[str, Any]], df_source: pd.DataFrame) -> dict[str, Any]:\n \"\"\"Build detailed column metadata.\"\"\"\n metadata: dict[str, Any] = {\n \"total_columns\": len(df_source.columns),\n \"mapped_columns\": len(config_list),\n \"unmapped_columns\": len(df_source.columns) - len(config_list),\n \"columns\": [],\n \"summary\": {\"vectorized_columns\": [], \"identifier_columns\": []},\n }\n\n for config in config_list:\n col_name = config.get(\"column_name\")\n vectorize = config.get(\"vectorize\") == \"True\" or config.get(\"vectorize\") is True\n identifier = config.get(\"identifier\") == \"True\" or config.get(\"identifier\") is True\n\n # Add to columns list\n metadata[\"columns\"].append(\n {\n \"name\": col_name,\n \"vectorize\": vectorize,\n \"identifier\": identifier,\n }\n )\n\n # Update summary\n if vectorize:\n metadata[\"summary\"][\"vectorized_columns\"].append(col_name)\n if identifier:\n metadata[\"summary\"][\"identifier_columns\"].append(col_name)\n\n return metadata\n\n async def _create_vector_store(\n self,\n df_source: pd.DataFrame,\n config_list: list[dict[str, Any]],\n embedding_function,\n ) -> Chroma:\n \"\"\"Create vector store following Local DB component pattern.\n\n Returns the Chroma instance so callers can use it for metrics updates.\n \"\"\"\n # Set up vector store directory\n vector_store_dir = await self._kb_path()\n if not vector_store_dir:\n msg = \"Knowledge base path is not set. Please create a new knowledge base first.\"\n raise ValueError(msg)\n vector_store_dir.mkdir(parents=True, exist_ok=True)\n\n # Convert DataFrame to Data objects (following Local DB pattern)\n data_objects = await self._convert_df_to_data_objects(df_source, config_list)\n\n # Create vector store\n chroma = Chroma(\n persist_directory=str(vector_store_dir),\n embedding_function=embedding_function,\n collection_name=self.knowledge_base,\n )\n\n # Convert Data objects to LangChain Documents\n documents = []\n for data_obj in data_objects:\n doc = data_obj.to_lc_document()\n documents.append(doc)\n\n # Add documents to vector store\n if documents:\n chroma.add_documents(documents)\n self.log(f\"Added {len(documents)} documents to vector store '{self.knowledge_base}'\")\n\n return chroma\n\n async def _convert_df_to_data_objects(\n self, df_source: pd.DataFrame, config_list: list[dict[str, Any]]\n ) -> list[Data]:\n \"\"\"Convert DataFrame to Data objects for vector store.\"\"\"\n data_objects: list[Data] = []\n\n # Set up vector store directory\n kb_path = await self._kb_path()\n\n # If we don't allow duplicates, we need to get the existing hashes\n chroma = Chroma(\n persist_directory=str(kb_path),\n collection_name=self.knowledge_base,\n )\n\n # Get all documents and their metadata\n all_docs = chroma.get()\n\n # Extract all _id values from metadata\n id_list = [metadata.get(\"_id\") for metadata in all_docs[\"metadatas\"] if metadata.get(\"_id\")]\n\n # Get column roles\n content_cols = []\n identifier_cols = []\n\n for config in config_list:\n col_name = config.get(\"column_name\")\n vectorize = config.get(\"vectorize\") == \"True\" or config.get(\"vectorize\") is True\n identifier = config.get(\"identifier\") == \"True\" or config.get(\"identifier\") is True\n\n if vectorize:\n content_cols.append(col_name)\n elif identifier:\n identifier_cols.append(col_name)\n\n # Convert each row to a Data object\n for _, row in df_source.iterrows():\n # Build content text from identifier columns using list comprehension\n identifier_parts = [str(row[col]) for col in content_cols if col in row and self._scalar_notna(row[col])]\n\n # Join all parts into a single string\n page_content = \" \".join(identifier_parts)\n\n # Build metadata from NON-vectorized columns only (simple key-value pairs)\n data_dict = {\n \"text\": page_content, # Main content for vectorization\n }\n\n # Add identifier columns if they exist\n if identifier_cols:\n identifier_parts = [\n str(row[col]) for col in identifier_cols if col in row and self._scalar_notna(row[col])\n ]\n page_content = \" \".join(identifier_parts)\n\n # Add metadata columns as simple key-value pairs\n for col in df_source.columns:\n if col not in content_cols and col in row and self._scalar_notna(row[col]):\n # Convert to simple types for Chroma metadata\n value = row[col]\n data_dict[col] = str(value) # Convert complex types to string\n\n # Hash the page_content for unique ID\n page_content_hash = hashlib.sha256(page_content.encode()).hexdigest()\n data_dict[\"_id\"] = page_content_hash\n\n # If duplicates are disallowed, and hash exists, prevent adding this row\n if not self.allow_duplicates and page_content_hash in id_list:\n self.log(f\"Skipping duplicate row with hash {page_content_hash}\")\n continue\n\n # Create Data object - everything except \"text\" becomes metadata\n data_obj = Data(data=data_dict)\n data_objects.append(data_obj)\n\n return data_objects\n\n def is_valid_collection_name(self, name, min_length: int = 3, max_length: int = 63) -> bool:\n \"\"\"Validates collection name against conditions 1-3.\n\n 1. Contains 3-63 characters\n 2. Starts and ends with alphanumeric character\n 3. Contains only alphanumeric characters, underscores, or hyphens.\n\n Args:\n name (str): Collection name to validate\n min_length (int): Minimum length of the name\n max_length (int): Maximum length of the name\n\n Returns:\n bool: True if valid, False otherwise\n \"\"\"\n # Check length (condition 1)\n if not (min_length <= len(name) <= max_length):\n return False\n\n # Check start/end with alphanumeric (condition 2)\n if not (name[0].isalnum() and name[-1].isalnum()):\n return False\n\n # Check allowed characters (condition 3)\n return re.match(r\"^[a-zA-Z0-9_-]+$\", name) is not None\n\n async def _kb_path(self) -> Path | None:\n # Check if we already have the path cached\n cached_path = getattr(self, \"_cached_kb_path\", None)\n if cached_path is not None:\n return cached_path\n\n # If not cached, compute it\n async with session_scope() as db:\n if not self.user_id:\n msg = \"User ID is required for fetching knowledge base path.\"\n raise ValueError(msg)\n current_user = await get_user_by_id(db, self.user_id)\n if not current_user:\n msg = f\"User with ID {self.user_id} not found.\"\n raise ValueError(msg)\n kb_user = current_user.username\n\n kb_root = self._get_kb_root()\n\n # Cache the result\n self._cached_kb_path = kb_root / kb_user / self.knowledge_base\n\n return self._cached_kb_path\n\n # ---------------------------------------------------------------------\n # OUTPUT METHODS\n # ---------------------------------------------------------------------\n async def build_kb_info(self) -> Data:\n \"\"\"Main ingestion routine → returns a dict with KB metadata.\"\"\"\n # Check if we're in Astra cloud environment and raise an error if we are.\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n try:\n input_value = self.input_df[0] if isinstance(self.input_df, list) else self.input_df\n df_source: DataFrame = convert_to_dataframe(input_value, auto_parse=False)\n\n # Validate column configuration (using Structured Output patterns)\n config_list = self._validate_column_config(df_source)\n column_metadata = self._build_column_metadata(config_list, df_source)\n\n # Read the embedding info from the knowledge base folder\n kb_path = await self._kb_path()\n if not kb_path:\n msg = \"Knowledge base path is not set. Please create a new knowledge base first.\"\n raise ValueError(msg)\n metadata_path = kb_path / \"embedding_metadata.json\"\n api_key = None\n model_selection = None\n\n # Read stored metadata\n if metadata_path.exists():\n settings_service = get_settings_service()\n stored_metadata = json.loads(metadata_path.read_text())\n\n # Prefer stored model_selection dict (new format)\n model_selection = stored_metadata.get(\"model_selection\")\n if model_selection:\n model_selection = [model_selection] if isinstance(model_selection, dict) else model_selection\n else:\n # Backward compat: reconstruct from old string-based metadata\n embedding_model_name = stored_metadata.get(\"embedding_model\")\n embedding_provider = stored_metadata.get(\"embedding_provider\", \"Unknown\")\n if embedding_model_name:\n # Look up full model info from available options\n try:\n all_options = get_embedding_model_options(user_id=self.user_id)\n match = next(\n (o for o in all_options if o.get(\"name\") == embedding_model_name),\n None,\n )\n if match:\n model_selection = [match]\n else:\n self.log(\n f\"Embedding model '{embedding_model_name}' (provider: {embedding_provider}) \"\n \"from stored metadata is no longer available in the model registry. \"\n \"Please re-create this knowledge base with a supported embedding model.\"\n )\n msg = (\n f\"Embedding model '{embedding_model_name}' is no longer recognized. \"\n \"The knowledge base was created with an older format and the model \"\n \"is not available in the current registry. \"\n \"Please re-create the knowledge base with a supported embedding model.\"\n )\n raise ValueError(msg)\n except ValueError:\n raise\n except Exception: # noqa: BLE001\n self.log(\n f\"Failed to look up embedding model '{embedding_model_name}' in registry. \"\n \"Please re-create this knowledge base with a supported embedding model.\"\n )\n msg = (\n f\"Could not look up embedding model '{embedding_model_name}' \"\n f\"(provider: {embedding_provider}). \"\n \"Please re-create the knowledge base with a supported embedding model.\"\n )\n raise ValueError(msg) # noqa: B904\n\n # Decrypt stored API key\n encrypted_key = stored_metadata.get(\"api_key\")\n if encrypted_key:\n try:\n api_key = decrypt_api_key(encrypted_key, settings_service)\n except (InvalidToken, TypeError, ValueError) as e:\n self.log(f\"Could not decrypt API key. Please provide it manually. Error: {e}\")\n\n # Check if a custom API key was provided\n if self.api_key:\n api_key = self.api_key\n if model_selection:\n self._save_embedding_metadata(\n kb_path=kb_path,\n model_selection=model_selection,\n api_key=api_key,\n )\n\n if not model_selection:\n msg = \"No embedding model configuration found. Please create the knowledge base first.\"\n raise ValueError(msg)\n\n # Build the embedding function via the shared utility\n embedding_function = get_embeddings(\n model=model_selection,\n user_id=self.user_id,\n api_key=api_key,\n chunk_size=self.chunk_size,\n )\n\n # Create vector store following Local DB component pattern\n chroma = await self._create_vector_store(df_source, config_list, embedding_function=embedding_function)\n\n # Save KB files (using File Component storage patterns)\n self._save_kb_files(kb_path, config_list)\n\n # Update embedding_metadata.json with accurate text metrics\n # so the KB modal and API show correct chunks/words/characters\n self._update_metadata_metrics(kb_path, chroma)\n\n # Build metadata response\n meta: dict[str, Any] = {\n \"kb_id\": str(uuid.uuid4()),\n \"kb_name\": self.knowledge_base,\n \"rows\": len(df_source),\n \"column_metadata\": column_metadata,\n \"path\": str(kb_path),\n \"config_columns\": len(config_list),\n \"timestamp\": datetime.now(tz=timezone.utc).isoformat(),\n }\n\n # Set status message\n self.status = f\"✅ KB **{self.knowledge_base}** saved · {len(df_source)} chunks.\"\n\n return Data(data=meta)\n\n except (OSError, ValueError, RuntimeError, KeyError) as e:\n msg = f\"Error during KB ingestion: {e}\"\n raise RuntimeError(msg) from e\n\n async def update_build_config(\n self,\n build_config,\n field_value: Any,\n field_name: str | None = None,\n ):\n \"\"\"Update build configuration based on provider selection.\"\"\"\n # Check if we're in Astra cloud environment and raise an error if we are.\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n\n # Populate the dialog's embedding model options so the ModelInput renders correctly\n try:\n dialog_template = (\n build_config[\"knowledge_base\"]\n .get(\"dialog_inputs\", {})\n .get(\"fields\", {})\n .get(\"data\", {})\n .get(\"node\", {})\n .get(\"template\", {})\n )\n if \"02_embedding_model\" in dialog_template:\n embedding_options = get_embedding_model_options(user_id=self.user_id)\n dialog_template[\"02_embedding_model\"][\"options\"] = embedding_options\n except Exception: # noqa: BLE001\n self.log(\"Failed to populate embedding model options in dialog\")\n\n # Create a new knowledge base\n if field_name == \"knowledge_base\":\n async with session_scope() as db:\n if not self.user_id:\n msg = \"User ID is required for fetching knowledge base list.\"\n raise ValueError(msg)\n current_user = await get_user_by_id(db, self.user_id)\n if not current_user:\n msg = f\"User with ID {self.user_id} not found.\"\n raise ValueError(msg)\n kb_user = current_user.username\n if isinstance(field_value, dict) and \"01_new_kb_name\" in field_value:\n # Validate the knowledge base name - Make sure it follows these rules:\n if not self.is_valid_collection_name(field_value[\"01_new_kb_name\"]):\n msg = f\"Invalid knowledge base name: {field_value['01_new_kb_name']}\"\n raise ValueError(msg)\n\n # The model selection comes from ModelInput as a list of dicts\n model_selection = field_value[\"02_embedding_model\"]\n if isinstance(model_selection, dict):\n model_selection = [model_selection]\n\n api_key = field_value.get(\"03_api_key\") or None\n\n # Build and validate the embedding model via the shared utility\n embed_model = get_embeddings(\n model=model_selection,\n user_id=self.user_id,\n api_key=api_key,\n )\n\n # Try to generate a dummy embedding to validate without blocking the event loop\n try:\n await asyncio.wait_for(\n asyncio.to_thread(embed_model.embed_query, \"test\"),\n timeout=10,\n )\n except TimeoutError as e:\n msg = \"Embedding validation timed out. Please verify network connectivity and key.\"\n raise ValueError(msg) from e\n except Exception as e:\n msg = f\"Embedding validation failed: {e!s}\"\n raise ValueError(msg) from e\n\n # Create the new knowledge base directory\n kb_path = _get_knowledge_bases_root_path() / kb_user / field_value[\"01_new_kb_name\"]\n kb_path.mkdir(parents=True, exist_ok=True)\n\n # Save the embedding metadata\n build_config[\"knowledge_base\"][\"value\"] = field_value[\"01_new_kb_name\"]\n self._save_embedding_metadata(\n kb_path=kb_path,\n model_selection=model_selection,\n api_key=api_key,\n )\n\n # Update the knowledge base options dynamically\n build_config[\"knowledge_base\"][\"options\"] = await get_knowledge_bases(\n _get_knowledge_bases_root_path(),\n user_id=self.user_id,\n )\n\n # If the selected knowledge base is not available, reset it\n if build_config[\"knowledge_base\"][\"value\"] not in build_config[\"knowledge_base\"][\"options\"]:\n build_config[\"knowledge_base\"][\"value\"] = None\n\n return build_config\n" + "value": "from __future__ import annotations\n\nimport asyncio\nimport hashlib\nimport json\nimport re\nimport uuid\nfrom dataclasses import asdict, dataclass, field\nfrom datetime import datetime, timezone\nfrom pathlib import Path\nfrom typing import TYPE_CHECKING, Any\n\nimport pandas as pd\nfrom cryptography.fernet import InvalidToken\nfrom langchain_chroma import Chroma\nfrom langflow.services.auth.utils import decrypt_api_key, encrypt_api_key\nfrom langflow.services.database.models.user.crud import get_user_by_id\n\nfrom lfx.base.knowledge_bases.knowledge_base_utils import get_knowledge_bases\nfrom lfx.base.models.unified_models import get_embedding_model_options, get_embeddings\nfrom lfx.components.processing.converter import convert_to_dataframe\nfrom lfx.custom import Component\nfrom lfx.io import (\n BoolInput,\n DropdownInput,\n HandleInput,\n IntInput,\n ModelInput,\n Output,\n SecretStrInput,\n StrInput,\n TableInput,\n)\nfrom lfx.schema.data import Data\nfrom lfx.schema.table import EditMode\nfrom lfx.services.deps import (\n get_settings_service,\n session_scope,\n)\nfrom lfx.utils.validate_cloud import raise_error_if_astra_cloud_disable_component\n\nif TYPE_CHECKING:\n from lfx.schema.dataframe import DataFrame\n\n_KNOWLEDGE_BASES_ROOT_PATH: Path | None = None\n\n# Error message to raise if we're in Astra cloud environment and the component is not supported.\nastra_error_msg = \"Knowledge ingestion is not supported in Astra cloud environment.\"\n\n\ndef _get_knowledge_bases_root_path() -> Path:\n \"\"\"Lazy load the knowledge bases root path from settings.\"\"\"\n global _KNOWLEDGE_BASES_ROOT_PATH # noqa: PLW0603\n if _KNOWLEDGE_BASES_ROOT_PATH is None:\n settings = get_settings_service().settings\n knowledge_directory = settings.knowledge_bases_dir\n if not knowledge_directory:\n msg = \"Knowledge bases directory is not set in the settings.\"\n raise ValueError(msg)\n _KNOWLEDGE_BASES_ROOT_PATH = Path(knowledge_directory).expanduser()\n return _KNOWLEDGE_BASES_ROOT_PATH\n\n\nclass KnowledgeIngestionComponent(Component):\n \"\"\"Create or append to Langflow Knowledge from a DataFrame.\"\"\"\n\n # ------ UI metadata ---------------------------------------------------\n display_name = \"Knowledge Ingestion\"\n description = \"Create or update knowledge in Langflow.\"\n icon = \"upload\"\n name = \"KnowledgeIngestion\"\n\n def __init__(self, *args, **kwargs) -> None:\n super().__init__(*args, **kwargs)\n self._cached_kb_path: Path | None = None\n\n @dataclass\n class NewKnowledgeBaseInput:\n functionality: str = \"create\"\n fields: dict[str, dict] = field(\n default_factory=lambda: {\n \"data\": {\n \"node\": {\n \"name\": \"create_knowledge_base\",\n \"description\": \"Create new knowledge in Langflow.\",\n \"display_name\": \"Create new knowledge\",\n \"field_order\": [\n \"01_new_kb_name\",\n \"02_embedding_model\",\n ],\n \"template\": {\n \"01_new_kb_name\": StrInput(\n name=\"new_kb_name\",\n display_name=\"Knowledge Name\",\n info=\"Name of the new knowledge to create.\",\n required=True,\n ),\n \"02_embedding_model\": ModelInput(\n name=\"embedding_model\",\n display_name=\"Choose Embedding Model\",\n info=(\n \"Select the embedding model to use for this knowledge base. \"\n \"Langflow uses the configured credentials for that model provider.\"\n ),\n required=True,\n model_type=\"embedding\",\n ),\n },\n },\n }\n }\n )\n\n # ------ Inputs --------------------------------------------------------\n inputs = [\n DropdownInput(\n name=\"knowledge_base\",\n display_name=\"Knowledge\",\n info=\"Select the knowledge to load data from.\",\n required=True,\n options=[],\n refresh_button=True,\n real_time_refresh=True,\n dialog_inputs=asdict(NewKnowledgeBaseInput()),\n ),\n HandleInput(\n name=\"input_df\",\n display_name=\"Input\",\n info=(\n \"Table with all original columns (already chunked / processed). \"\n \"Accepts Message, Data, or DataFrame. If Message or Data is provided, \"\n \"it is converted to a DataFrame automatically.\"\n ),\n input_types=[\"Message\", \"Data\", \"JSON\", \"DataFrame\", \"Table\"],\n required=True,\n ),\n TableInput(\n name=\"column_config\",\n display_name=\"Column Configuration\",\n info=\"Configure column behavior for the knowledge base.\",\n required=True,\n table_schema=[\n {\n \"name\": \"column_name\",\n \"display_name\": \"Column Name\",\n \"type\": \"str\",\n \"description\": \"Name of the column in the source DataFrame\",\n \"edit_mode\": EditMode.INLINE,\n },\n {\n \"name\": \"vectorize\",\n \"display_name\": \"Vectorize\",\n \"type\": \"boolean\",\n \"description\": \"Create embeddings for this column\",\n \"default\": False,\n \"edit_mode\": EditMode.INLINE,\n },\n {\n \"name\": \"identifier\",\n \"display_name\": \"Identifier\",\n \"type\": \"boolean\",\n \"description\": \"Use this column as unique identifier\",\n \"default\": False,\n \"edit_mode\": EditMode.INLINE,\n },\n ],\n value=[\n {\n \"column_name\": \"text\",\n \"vectorize\": True,\n \"identifier\": True,\n },\n ],\n ),\n IntInput(\n name=\"chunk_size\",\n display_name=\"Chunk Size\",\n info=\"Batch size for processing embeddings\",\n advanced=True,\n value=1000,\n ),\n SecretStrInput(\n name=\"api_key\",\n display_name=\"Embedding Provider API Key\",\n info=\"Overrides global provider settings. Leave blank to use your pre-configured API Key.\",\n advanced=True,\n required=False,\n ),\n BoolInput(\n name=\"allow_duplicates\",\n display_name=\"Allow Duplicates\",\n info=\"Allow duplicate rows in the knowledge base\",\n advanced=True,\n value=False,\n ),\n ]\n\n # ------ Outputs -------------------------------------------------------\n outputs = [Output(display_name=\"Results\", name=\"dataframe_output\", method=\"build_kb_info\")]\n\n # ------ Internal helpers ---------------------------------------------\n def _get_kb_root(self) -> Path:\n \"\"\"Return the root directory for knowledge bases.\"\"\"\n return _get_knowledge_bases_root_path()\n\n @staticmethod\n def _scalar_notna(value) -> bool:\n \"\"\"Check if a value is not NA, safely handling arrays and sequences.\n\n ``pd.notna`` returns an array when given an array-like input, which\n cannot be used directly in a boolean context. This helper collapses the\n result to a single scalar ``bool``.\n \"\"\"\n result = pd.notna(value)\n # If result is array-like (numpy array, list, etc.), treat non-empty arrays as \"present\"\n if hasattr(result, \"__iter__\") and not isinstance(result, str):\n import numpy as np\n\n arr = np.asarray(result)\n return arr.size > 0 and arr.all()\n return bool(result)\n\n def _validate_column_config(self, df_source: pd.DataFrame) -> list[dict[str, Any]]:\n \"\"\"Validate column configuration using Structured Output patterns.\"\"\"\n if not self.column_config:\n msg = \"Column configuration cannot be empty\"\n raise ValueError(msg)\n\n # Convert table input to list of dicts (similar to Structured Output)\n config_list = self.column_config if isinstance(self.column_config, list) else []\n\n # Validate column names exist in DataFrame\n df_columns = set(df_source.columns)\n for config in config_list:\n col_name = config.get(\"column_name\")\n if col_name not in df_columns:\n msg = f\"Column '{col_name}' not found in DataFrame. Available columns: {sorted(df_columns)}\"\n raise ValueError(msg)\n\n return config_list\n\n def _build_embedding_metadata(\n self,\n model_selection: list[dict[str, Any]],\n api_key: str | None = None,\n ) -> dict[str, Any]:\n \"\"\"Build embedding model metadata from a model selection dict.\n\n Args:\n model_selection: Model selection list from ModelInput\n (e.g. [{'name': ..., 'provider': ..., 'metadata': ...}])\n api_key: Optional runtime API key override.\n \"\"\"\n model_dict = model_selection[0] if isinstance(model_selection, list) else model_selection\n embedding_model = model_dict.get(\"name\", \"\")\n embedding_provider = model_dict.get(\"provider\", \"Unknown\")\n\n api_key_to_save = None\n if api_key and hasattr(api_key, \"get_secret_value\"):\n api_key_to_save = api_key.get_secret_value()\n elif isinstance(api_key, str):\n api_key_to_save = api_key\n\n encrypted_api_key = None\n if api_key_to_save:\n settings_service = get_settings_service()\n try:\n encrypted_api_key = encrypt_api_key(api_key_to_save, settings_service=settings_service)\n except (TypeError, ValueError) as e:\n self.log(f\"Could not encrypt API key: {e}\")\n\n return {\n \"embedding_provider\": embedding_provider,\n \"embedding_model\": embedding_model,\n \"model_selection\": model_dict, # Store full selection for get_embeddings() reconstruction\n \"api_key\": encrypted_api_key,\n \"api_key_used\": bool(api_key),\n \"chunk_size\": self.chunk_size,\n \"created_at\": datetime.now(timezone.utc).isoformat(),\n }\n\n def _save_embedding_metadata(\n self,\n kb_path: Path,\n model_selection: list[dict[str, Any]],\n api_key: str | None = None,\n ) -> None:\n \"\"\"Save embedding model metadata.\"\"\"\n embedding_metadata = self._build_embedding_metadata(model_selection, api_key)\n metadata_path = kb_path / \"embedding_metadata.json\"\n metadata_path.write_text(json.dumps(embedding_metadata, indent=2))\n\n def _update_metadata_metrics(self, kb_path: Path, chroma: Chroma) -> None:\n \"\"\"Update embedding_metadata.json with accurate chunk/word/character counts.\n\n This ensures the Knowledge Base modal displays correct stats after\n component-based ingestion, matching the behavior of API-based ingestion.\n Delegates to KBAnalysisHelper.update_text_metrics to avoid duplicating\n the batched metrics counting logic.\n \"\"\"\n import chromadb.errors\n from langflow.api.utils.kb_helpers import KBAnalysisHelper, KBStorageHelper\n\n metadata_path = kb_path / \"embedding_metadata.json\"\n if not metadata_path.exists():\n return\n\n try:\n metadata = json.loads(metadata_path.read_text())\n KBAnalysisHelper.update_text_metrics(kb_path, metadata, chroma)\n metadata[\"size\"] = KBStorageHelper.get_directory_size(kb_path)\n metadata_path.write_text(json.dumps(metadata, indent=2))\n except (OSError, ValueError, TypeError, json.JSONDecodeError, chromadb.errors.ChromaError) as e:\n self.log(f\"Warning: Could not update metadata metrics: {e}\")\n\n def _save_kb_files(\n self,\n kb_path: Path,\n config_list: list[dict[str, Any]],\n ) -> None:\n \"\"\"Save KB files using File Component storage patterns.\"\"\"\n try:\n # Create directory (following File Component patterns)\n kb_path.mkdir(parents=True, exist_ok=True)\n\n # Save column configuration\n # Only do this if the file doesn't exist already\n cfg_path = kb_path / \"schema.json\"\n if not cfg_path.exists():\n cfg_path.write_text(json.dumps(config_list, indent=2))\n\n except (OSError, TypeError, ValueError) as e:\n self.log(f\"Error saving KB files: {e}\")\n\n def _build_column_metadata(self, config_list: list[dict[str, Any]], df_source: pd.DataFrame) -> dict[str, Any]:\n \"\"\"Build detailed column metadata.\"\"\"\n metadata: dict[str, Any] = {\n \"total_columns\": len(df_source.columns),\n \"mapped_columns\": len(config_list),\n \"unmapped_columns\": len(df_source.columns) - len(config_list),\n \"columns\": [],\n \"summary\": {\"vectorized_columns\": [], \"identifier_columns\": []},\n }\n\n for config in config_list:\n col_name = config.get(\"column_name\")\n vectorize = config.get(\"vectorize\") == \"True\" or config.get(\"vectorize\") is True\n identifier = config.get(\"identifier\") == \"True\" or config.get(\"identifier\") is True\n\n # Add to columns list\n metadata[\"columns\"].append(\n {\n \"name\": col_name,\n \"vectorize\": vectorize,\n \"identifier\": identifier,\n }\n )\n\n # Update summary\n if vectorize:\n metadata[\"summary\"][\"vectorized_columns\"].append(col_name)\n if identifier:\n metadata[\"summary\"][\"identifier_columns\"].append(col_name)\n\n return metadata\n\n async def _create_vector_store(\n self,\n df_source: pd.DataFrame,\n config_list: list[dict[str, Any]],\n embedding_function,\n ) -> Chroma:\n \"\"\"Create vector store following Local DB component pattern.\n\n Returns the Chroma instance so callers can use it for metrics updates.\n \"\"\"\n # Set up vector store directory\n vector_store_dir = await self._kb_path()\n if not vector_store_dir:\n msg = \"Knowledge base path is not set. Please create a new knowledge base first.\"\n raise ValueError(msg)\n vector_store_dir.mkdir(parents=True, exist_ok=True)\n\n # Convert DataFrame to Data objects (following Local DB pattern)\n data_objects = await self._convert_df_to_data_objects(df_source, config_list)\n\n # Create vector store\n chroma = Chroma(\n persist_directory=str(vector_store_dir),\n embedding_function=embedding_function,\n collection_name=self.knowledge_base,\n )\n\n # Convert Data objects to LangChain Documents\n documents = []\n for data_obj in data_objects:\n doc = data_obj.to_lc_document()\n documents.append(doc)\n\n # Add documents to vector store\n if documents:\n chroma.add_documents(documents)\n self.log(f\"Added {len(documents)} documents to vector store '{self.knowledge_base}'\")\n\n return chroma\n\n async def _convert_df_to_data_objects(\n self, df_source: pd.DataFrame, config_list: list[dict[str, Any]]\n ) -> list[Data]:\n \"\"\"Convert DataFrame to Data objects for vector store.\"\"\"\n data_objects: list[Data] = []\n\n # Set up vector store directory\n kb_path = await self._kb_path()\n\n # If we don't allow duplicates, we need to get the existing hashes\n chroma = Chroma(\n persist_directory=str(kb_path),\n collection_name=self.knowledge_base,\n )\n\n # Get all documents and their metadata\n all_docs = chroma.get()\n\n # Extract all _id values from metadata\n id_list = [metadata.get(\"_id\") for metadata in all_docs[\"metadatas\"] if metadata.get(\"_id\")]\n\n # Get column roles\n content_cols = []\n identifier_cols = []\n\n for config in config_list:\n col_name = config.get(\"column_name\")\n vectorize = config.get(\"vectorize\") == \"True\" or config.get(\"vectorize\") is True\n identifier = config.get(\"identifier\") == \"True\" or config.get(\"identifier\") is True\n\n if vectorize:\n content_cols.append(col_name)\n elif identifier:\n identifier_cols.append(col_name)\n\n # Convert each row to a Data object\n for _, row in df_source.iterrows():\n # Build content text from identifier columns using list comprehension\n identifier_parts = [str(row[col]) for col in content_cols if col in row and self._scalar_notna(row[col])]\n\n # Join all parts into a single string\n page_content = \" \".join(identifier_parts)\n\n # Build metadata from NON-vectorized columns only (simple key-value pairs)\n data_dict = {\n \"text\": page_content, # Main content for vectorization\n }\n\n # Add identifier columns if they exist\n if identifier_cols:\n identifier_parts = [\n str(row[col]) for col in identifier_cols if col in row and self._scalar_notna(row[col])\n ]\n page_content = \" \".join(identifier_parts)\n\n # Add metadata columns as simple key-value pairs\n for col in df_source.columns:\n if col not in content_cols and col in row and self._scalar_notna(row[col]):\n # Convert to simple types for Chroma metadata\n value = row[col]\n data_dict[col] = str(value) # Convert complex types to string\n\n # Hash the page_content for unique ID\n page_content_hash = hashlib.sha256(page_content.encode()).hexdigest()\n data_dict[\"_id\"] = page_content_hash\n\n # If duplicates are disallowed, and hash exists, prevent adding this row\n if not self.allow_duplicates and page_content_hash in id_list:\n self.log(f\"Skipping duplicate row with hash {page_content_hash}\")\n continue\n\n # Create Data object - everything except \"text\" becomes metadata\n data_obj = Data(data=data_dict)\n data_objects.append(data_obj)\n\n return data_objects\n\n def is_valid_collection_name(self, name, min_length: int = 3, max_length: int = 63) -> bool:\n \"\"\"Validates collection name against conditions 1-3.\n\n 1. Contains 3-63 characters\n 2. Starts and ends with alphanumeric character\n 3. Contains only alphanumeric characters, underscores, or hyphens.\n\n Args:\n name (str): Collection name to validate\n min_length (int): Minimum length of the name\n max_length (int): Maximum length of the name\n\n Returns:\n bool: True if valid, False otherwise\n \"\"\"\n # Check length (condition 1)\n if not (min_length <= len(name) <= max_length):\n return False\n\n # Check start/end with alphanumeric (condition 2)\n if not (name[0].isalnum() and name[-1].isalnum()):\n return False\n\n # Check allowed characters (condition 3)\n return re.match(r\"^[a-zA-Z0-9_-]+$\", name) is not None\n\n async def _kb_path(self) -> Path | None:\n # Check if we already have the path cached\n cached_path = getattr(self, \"_cached_kb_path\", None)\n if cached_path is not None:\n return cached_path\n\n # If not cached, compute it\n async with session_scope() as db:\n if not self.user_id:\n msg = \"User ID is required for fetching knowledge base path.\"\n raise ValueError(msg)\n current_user = await get_user_by_id(db, self.user_id)\n if not current_user:\n msg = f\"User with ID {self.user_id} not found.\"\n raise ValueError(msg)\n kb_user = current_user.username\n\n kb_root = self._get_kb_root()\n\n # Cache the result\n self._cached_kb_path = kb_root / kb_user / self.knowledge_base\n\n return self._cached_kb_path\n\n # ---------------------------------------------------------------------\n # OUTPUT METHODS\n # ---------------------------------------------------------------------\n async def build_kb_info(self) -> Data:\n \"\"\"Main ingestion routine → returns a dict with KB metadata.\"\"\"\n # Check if we're in Astra cloud environment and raise an error if we are.\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n try:\n input_value = self.input_df[0] if isinstance(self.input_df, list) else self.input_df\n df_source: DataFrame = convert_to_dataframe(input_value, auto_parse=False)\n\n # Validate column configuration (using Structured Output patterns)\n config_list = self._validate_column_config(df_source)\n column_metadata = self._build_column_metadata(config_list, df_source)\n\n # Read the embedding info from the knowledge base folder\n kb_path = await self._kb_path()\n if not kb_path:\n msg = \"Knowledge base path is not set. Please create a new knowledge base first.\"\n raise ValueError(msg)\n metadata_path = kb_path / \"embedding_metadata.json\"\n api_key = None\n model_selection = None\n\n # Read stored metadata\n if metadata_path.exists():\n settings_service = get_settings_service()\n stored_metadata = json.loads(metadata_path.read_text())\n\n # Prefer stored model_selection dict (new format)\n model_selection = stored_metadata.get(\"model_selection\")\n if model_selection:\n model_selection = [model_selection] if isinstance(model_selection, dict) else model_selection\n else:\n # Backward compat: reconstruct from old string-based metadata\n embedding_model_name = stored_metadata.get(\"embedding_model\")\n embedding_provider = stored_metadata.get(\"embedding_provider\", \"Unknown\")\n if embedding_model_name:\n # Look up full model info from available options\n try:\n all_options = get_embedding_model_options(user_id=self.user_id)\n match = next(\n (o for o in all_options if o.get(\"name\") == embedding_model_name),\n None,\n )\n if match:\n model_selection = [match]\n else:\n self.log(\n f\"Embedding model '{embedding_model_name}' (provider: {embedding_provider}) \"\n \"from stored metadata is no longer available in the model registry. \"\n \"Please re-create this knowledge base with a supported embedding model.\"\n )\n msg = (\n f\"Embedding model '{embedding_model_name}' is no longer recognized. \"\n \"The knowledge base was created with an older format and the model \"\n \"is not available in the current registry. \"\n \"Please re-create the knowledge base with a supported embedding model.\"\n )\n raise ValueError(msg)\n except ValueError:\n raise\n except Exception: # noqa: BLE001\n self.log(\n f\"Failed to look up embedding model '{embedding_model_name}' in registry. \"\n \"Please re-create this knowledge base with a supported embedding model.\"\n )\n msg = (\n f\"Could not look up embedding model '{embedding_model_name}' \"\n f\"(provider: {embedding_provider}). \"\n \"Please re-create the knowledge base with a supported embedding model.\"\n )\n raise ValueError(msg) # noqa: B904\n\n # Decrypt stored API key\n encrypted_key = stored_metadata.get(\"api_key\")\n if encrypted_key:\n try:\n api_key = decrypt_api_key(encrypted_key, settings_service)\n except (InvalidToken, TypeError, ValueError) as e:\n self.log(f\"Could not decrypt API key. Please provide it manually. Error: {e}\")\n\n # Check if a custom API key was provided\n if self.api_key:\n api_key = self.api_key\n if model_selection:\n self._save_embedding_metadata(\n kb_path=kb_path,\n model_selection=model_selection,\n api_key=api_key,\n )\n\n if not model_selection:\n msg = \"No embedding model configuration found. Please create the knowledge base first.\"\n raise ValueError(msg)\n\n # Build the embedding function via the shared utility\n embedding_function = get_embeddings(\n model=model_selection,\n user_id=self.user_id,\n api_key=api_key,\n chunk_size=self.chunk_size,\n )\n\n # Create vector store following Local DB component pattern\n chroma = await self._create_vector_store(df_source, config_list, embedding_function=embedding_function)\n\n # Save KB files (using File Component storage patterns)\n self._save_kb_files(kb_path, config_list)\n\n # Update embedding_metadata.json with accurate text metrics\n # so the KB modal and API show correct chunks/words/characters\n self._update_metadata_metrics(kb_path, chroma)\n\n # Build metadata response\n meta: dict[str, Any] = {\n \"kb_id\": str(uuid.uuid4()),\n \"kb_name\": self.knowledge_base,\n \"rows\": len(df_source),\n \"column_metadata\": column_metadata,\n \"path\": str(kb_path),\n \"config_columns\": len(config_list),\n \"timestamp\": datetime.now(tz=timezone.utc).isoformat(),\n }\n\n # Set status message\n self.status = f\"✅ KB **{self.knowledge_base}** saved · {len(df_source)} chunks.\"\n\n return Data(data=meta)\n\n except (OSError, ValueError, RuntimeError, KeyError) as e:\n msg = f\"Error during KB ingestion: {e}\"\n raise RuntimeError(msg) from e\n\n async def update_build_config(\n self,\n build_config,\n field_value: Any,\n field_name: str | None = None,\n ):\n \"\"\"Update build configuration based on provider selection.\"\"\"\n # Check if we're in Astra cloud environment and raise an error if we are.\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n\n # Populate the dialog's embedding model options so the ModelInput renders correctly\n try:\n dialog_template = (\n build_config[\"knowledge_base\"]\n .get(\"dialog_inputs\", {})\n .get(\"fields\", {})\n .get(\"data\", {})\n .get(\"node\", {})\n .get(\"template\", {})\n )\n if \"02_embedding_model\" in dialog_template:\n embedding_options = get_embedding_model_options(user_id=self.user_id)\n dialog_template[\"02_embedding_model\"][\"options\"] = embedding_options\n except Exception: # noqa: BLE001\n self.log(\"Failed to populate embedding model options in dialog\")\n\n # Create a new knowledge base\n if field_name == \"knowledge_base\":\n async with session_scope() as db:\n if not self.user_id:\n msg = \"User ID is required for fetching knowledge base list.\"\n raise ValueError(msg)\n current_user = await get_user_by_id(db, self.user_id)\n if not current_user:\n msg = f\"User with ID {self.user_id} not found.\"\n raise ValueError(msg)\n kb_user = current_user.username\n if isinstance(field_value, dict) and \"01_new_kb_name\" in field_value:\n # Validate the knowledge base name - Make sure it follows these rules:\n if not self.is_valid_collection_name(field_value[\"01_new_kb_name\"]):\n msg = f\"Invalid knowledge base name: {field_value['01_new_kb_name']}\"\n raise ValueError(msg)\n\n # The model selection comes from ModelInput as a list of dicts\n model_selection = field_value[\"02_embedding_model\"]\n if isinstance(model_selection, dict):\n model_selection = [model_selection]\n\n # Build and validate the embedding model via the shared utility\n embed_model = get_embeddings(\n model=model_selection,\n user_id=self.user_id,\n )\n\n # Try to generate a dummy embedding to validate without blocking the event loop\n try:\n await asyncio.wait_for(\n asyncio.to_thread(embed_model.embed_query, \"test\"),\n timeout=10,\n )\n except TimeoutError as e:\n msg = \"Embedding validation timed out. Please verify network connectivity and key.\"\n raise ValueError(msg) from e\n except Exception as e:\n msg = f\"Embedding validation failed: {e!s}\"\n raise ValueError(msg) from e\n\n # Create the new knowledge base directory\n kb_path = _get_knowledge_bases_root_path() / kb_user / field_value[\"01_new_kb_name\"]\n kb_path.mkdir(parents=True, exist_ok=True)\n\n # Save the embedding metadata\n build_config[\"knowledge_base\"][\"value\"] = field_value[\"01_new_kb_name\"]\n self._save_embedding_metadata(\n kb_path=kb_path,\n model_selection=model_selection,\n )\n\n # Update the knowledge base options dynamically\n build_config[\"knowledge_base\"][\"options\"] = await get_knowledge_bases(\n _get_knowledge_bases_root_path(),\n user_id=self.user_id,\n )\n\n # If the selected knowledge base is not available, reset it\n if build_config[\"knowledge_base\"][\"value\"] not in build_config[\"knowledge_base\"][\"options\"]:\n build_config[\"knowledge_base\"][\"value\"] = None\n\n return build_config\n" }, "column_config": { "_input_type": "TableInput", @@ -69368,8 +69368,7 @@ "display_name": "Create new knowledge", "field_order": [ "01_new_kb_name", - "02_embedding_model", - "03_api_key" + "02_embedding_model" ], "name": "create_knowledge_base", "template": { @@ -69410,7 +69409,7 @@ } } }, - "info": "Select the embedding model to use for this knowledge base.", + "info": "Select the embedding model to use for this knowledge base. Langflow uses the configured credentials for that model provider.", "input_types": [ "Embeddings" ], @@ -69429,25 +69428,6 @@ "track_in_telemetry": false, "type": "model", "value": "" - }, - "03_api_key": { - "_input_type": "SecretStrInput", - "advanced": true, - "display_name": "Embedding Provider API Key", - "dynamic": false, - "info": "Optional API key override used to validate and save this knowledge base.", - "input_types": [], - "load_from_db": true, - "name": "api_key", - "override_skip": false, - "password": true, - "placeholder": "", - "required": false, - "show": true, - "title_case": false, - "track_in_telemetry": false, - "type": "str", - "value": "" } } } @@ -118124,6 +118104,6 @@ "num_components": 355, "num_modules": 97 }, - "sha256": "d4c0fe9459a9aa3ab17010e67cb69fd63d02e161b022ab1c949120eebaf59877", + "sha256": "dcd6d6fca00b87a3b8bb21c9c7f765ab2d19be06d1c0beeb39db07802f4d01c0", "version": "0.4.0" } diff --git a/src/lfx/src/lfx/components/files_and_knowledge/ingestion.py b/src/lfx/src/lfx/components/files_and_knowledge/ingestion.py index 442c075185..02f10438f8 100644 --- a/src/lfx/src/lfx/components/files_and_knowledge/ingestion.py +++ b/src/lfx/src/lfx/components/files_and_knowledge/ingestion.py @@ -87,7 +87,6 @@ class KnowledgeIngestionComponent(Component): "field_order": [ "01_new_kb_name", "02_embedding_model", - "03_api_key", ], "template": { "01_new_kb_name": StrInput( @@ -99,17 +98,13 @@ class KnowledgeIngestionComponent(Component): "02_embedding_model": ModelInput( name="embedding_model", display_name="Choose Embedding Model", - info="Select the embedding model to use for this knowledge base.", + info=( + "Select the embedding model to use for this knowledge base. " + "Langflow uses the configured credentials for that model provider." + ), required=True, model_type="embedding", ), - "03_api_key": SecretStrInput( - name="api_key", - display_name="Embedding Provider API Key", - info="Optional API key override used to validate and save this knowledge base.", - required=False, - advanced=True, - ), }, }, } @@ -254,7 +249,7 @@ class KnowledgeIngestionComponent(Component): Args: model_selection: Model selection list from ModelInput (e.g. [{'name': ..., 'provider': ..., 'metadata': ...}]) - api_key: Optional API key override. + api_key: Optional runtime API key override. """ model_dict = model_selection[0] if isinstance(model_selection, list) else model_selection embedding_model = model_dict.get("name", "") @@ -718,13 +713,10 @@ class KnowledgeIngestionComponent(Component): if isinstance(model_selection, dict): model_selection = [model_selection] - api_key = field_value.get("03_api_key") or None - # Build and validate the embedding model via the shared utility embed_model = get_embeddings( model=model_selection, user_id=self.user_id, - api_key=api_key, ) # Try to generate a dummy embedding to validate without blocking the event loop @@ -749,7 +741,6 @@ class KnowledgeIngestionComponent(Component): self._save_embedding_metadata( kb_path=kb_path, model_selection=model_selection, - api_key=api_key, ) # Update the knowledge base options dynamically