From b8f887bae6dc2dede443aae9eba3fda82859b2cb Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Thu, 25 Jun 2026 14:26:07 +0000 Subject: [PATCH] [autofix.ci] apply automated fixes --- src/lfx/src/lfx/_assets/component_index.json | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/lfx/src/lfx/_assets/component_index.json b/src/lfx/src/lfx/_assets/component_index.json index 77ea68264c..79d5a57bc3 100644 --- a/src/lfx/src/lfx/_assets/component_index.json +++ b/src/lfx/src/lfx/_assets/component_index.json @@ -7722,7 +7722,7 @@ "icon": "database", "legacy": false, "metadata": { - "code_hash": "26a843069b81", + "code_hash": "37863c267333", "dependencies": { "dependencies": [ { @@ -7868,7 +7868,7 @@ "show": true, "title_case": false, "type": "code", - "value": "\"\"\"Unified Knowledge component โ€” ingest into or retrieve from a knowledge base.\n\nThis component merges what used to live in ``ingestion.py`` and\n``retrieval.py`` into a single, mode-driven component. A ``TabInput``\n(\"๐Ÿ“ฅ Ingest\" / \"๐Ÿ” Retrieve\") drives which inputs and which output are\nvisible. The merged shape gives users one node per KB instead of two\nparallel ones that always shared the same ``knowledge_base`` picker,\nembedding-model metadata, and backend registry.\n\nBoth legacy classes (``KnowledgeIngestionComponent``,\n``KnowledgeBaseComponent``) remain importable as thin subclasses of\n``KnowledgeComponent`` โ€” see the sibling ``ingestion.py`` / ``retrieval.py``\nmodules โ€” so saved flows continue to load unchanged.\n\"\"\"\n\nfrom __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 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\n\nfrom lfx.base.knowledge_bases.backends import BackendType, BaseVectorStoreBackend, create_backend\nfrom lfx.base.knowledge_bases.ingestion_sources.base import (\n IngestionItemResult,\n IngestionItemStatus,\n IngestionRunStatus,\n IngestionSummary,\n)\nfrom lfx.base.knowledge_bases.ingestion_sources.flow_component import FlowComponentSource\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.base.vectorstores.chroma_security import chroma_langchain_collection_kwargs\nfrom lfx.components.files_and_knowledge._kb_paths import (\n get_knowledge_bases_root_path as _get_knowledge_bases_root_path,\n)\nfrom lfx.components.files_and_knowledge._kb_paths import (\n load_kb_metadata,\n)\nfrom lfx.components.processing.converter import convert_to_dataframe\nfrom lfx.custom import Component\nfrom lfx.io import (\n BoolInput,\n DBProviderInput,\n DropdownInput,\n HandleInput,\n IntInput,\n MessageTextInput,\n ModelInput,\n Output,\n SecretStrInput,\n StrInput,\n TabInput,\n TableInput,\n)\nfrom lfx.log.logger import logger\nfrom lfx.schema.data import Data\nfrom lfx.schema.dataframe import DataFrame\nfrom lfx.schema.dotdict import dotdict\nfrom lfx.schema.table import EditMode\nfrom lfx.services.deps import (\n get_settings_service,\n session_scope,\n)\nfrom lfx.utils.component_utils import set_current_fields, set_field_display\nfrom lfx.utils.validate_cloud import raise_error_if_astra_cloud_disable_component\n\n\ndef _inputs_for_mode(default_mode: str) -> list:\n \"\"\"Return a fresh copy of the canonical inputs list with show flags set for the given mode.\n\n Used by the legacy subclasses so a saved flow keyed on\n ``KnowledgeIngestion`` / ``KnowledgeBase`` lands on a node template\n whose default visibility already matches the pinned mode โ€” no\n \"wait for the user to click the mode tab\" UX gap on load.\n \"\"\"\n always_visible = set(KnowledgeComponent.default_keys) | set(KnowledgeComponent.mode_config[default_mode])\n inputs_copy = []\n for inp in KnowledgeComponent.inputs:\n clone = inp.model_copy(deep=True) if hasattr(inp, \"model_copy\") else inp\n if clone.name == \"mode\":\n clone.value = default_mode\n else:\n clone.show = clone.name in always_visible\n inputs_copy.append(clone)\n return inputs_copy\n\n\nif TYPE_CHECKING:\n from pathlib import Path\n\n# Mode constants. Plain-text labels (no emoji) for consistency with the rest of\n# Langflow's TabInput palette โ€” see ``MemoryComponent`` for the same convention.\nMODE_INGEST = \"Ingest\"\nMODE_RETRIEVE = \"Retrieve\"\n\n\ndef _is_retrieve_mode(value: Any) -> bool:\n \"\"\"Lenient mode check: treats any label containing 'Retrieve' as retrieve mode.\n\n Older saved flows may carry the emoji-prefixed labels (\"๐Ÿ“ฅ Ingest\" /\n \"๐Ÿ” Retrieve\") this component used to ship with; substring matching\n keeps those loading without forcing a flow rewrite.\n \"\"\"\n return isinstance(value, str) and \"Retrieve\" in value\n\n\n# Error message used by both the ingest and retrieve paths when the user is\n# running against an Astra cloud environment that disables these flows.\nastra_error_msg = \"Knowledge ingestion and retrieval are not supported in Astra cloud environment.\"\n\n_DEFAULT_OPENSEARCH_CONFIG = {\n \"url_variable\": \"OPENSEARCH_URL\",\n \"username_variable\": \"OPENSEARCH_USERNAME\",\n \"password_variable\": \"OPENSEARCH_PASSWORD\", # pragma: allowlist secret\n \"index_name\": \"\",\n \"vector_field\": \"vector_field\",\n \"text_field\": \"text\",\n}\n\n_DEFAULT_CHROMA_CLOUD_CONFIG = {\n \"mode\": \"cloud\",\n \"tenant_variable\": \"CHROMA_TENANT\",\n \"database_variable\": \"CHROMA_DATABASE\",\n \"api_key_variable\": \"CHROMA_API_KEY\", # pragma: allowlist secret\n}\n\n\nclass KnowledgeComponent(Component):\n \"\"\"One component for both writing into and reading from a Langflow knowledge base.\n\n A ``TabInput`` switches between ingestion and retrieval. The\n ``update_build_config`` / ``update_outputs`` hooks hide the inputs\n and the output that don't apply to the current mode so the canvas\n node stays focused.\n \"\"\"\n\n display_name = \"Knowledge\"\n description = \"Ingest into or retrieve from a Langflow knowledge base.\"\n icon = \"database\"\n name = \"Knowledge\"\n\n # ------ Mode โ†’ visible-fields wiring ---------------------------------\n # ``default_keys`` are inputs always visible regardless of mode.\n # ``mode_config`` lists the inputs unique to each mode; everything outside\n # this set is hidden when its mode is not selected.\n default_keys: list[str] = [\"mode\", \"knowledge_base\"]\n mode_config: dict[str, list[str]] = {\n MODE_INGEST: [\n \"input_df\",\n \"column_config\",\n \"chunk_size\",\n \"api_key\",\n \"allow_duplicates\",\n \"metadata_json\",\n ],\n MODE_RETRIEVE: [\n \"search_query\",\n \"top_k\",\n \"include_metadata\",\n \"include_embeddings\",\n \"metadata_filter\",\n ],\n }\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 Base\",\n \"field_order\": [\n \"01_new_kb_name\",\n \"02_embedding_model\",\n \"03_knowledge_backend\",\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 \"03_knowledge_backend\": DBProviderInput(\n name=\"knowledge_backend\",\n display_name=\"DB Provider\",\n info=(\n \"Select where this knowledge base stores vectors. \"\n \"OpenSearch uses the global DB Providers settings.\"\n ),\n required=True,\n ),\n },\n },\n }\n }\n )\n\n # ------ Inputs --------------------------------------------------------\n # ``knowledge_base`` is kept at position 0 so legacy tests that check\n # ``component.inputs[0].dialog_inputs`` continue to work.\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 TabInput(\n name=\"mode\",\n display_name=\"Mode\",\n options=[MODE_INGEST, MODE_RETRIEVE],\n value=MODE_INGEST,\n info=\"Switch between writing new data into the knowledge base and querying it.\",\n real_time_refresh=True,\n tool_mode=True,\n ),\n # --- Ingest-only inputs (default-shown; hidden when mode == Retrieve) -\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 dynamic=True,\n show=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 dynamic=True,\n show=True,\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 dynamic=True,\n show=True,\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 dynamic=True,\n show=True,\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 dynamic=True,\n show=True,\n ),\n StrInput(\n name=\"metadata_json\",\n display_name=\"Metadata\",\n info=(\n \"Optional JSON object of user metadata applied to every chunk produced by this \"\n 'run (e.g. {\"tag\": \"invoice\", \"year\": \"2026\"}). Same shape as the upload modal '\n \"Metadata section so chunks browser filters + Knowledge retrieval metadata_filter \"\n \"work uniformly across upload, folder, and flow-driven ingestion. Malformed JSON is \"\n \"ignored with a warning rather than failing the run.\"\n ),\n advanced=True,\n required=False,\n dynamic=True,\n show=True,\n ),\n # --- Retrieve-only inputs (default-hidden; shown when mode == Retrieve) -\n MessageTextInput(\n name=\"search_query\",\n display_name=\"Search Query\",\n info=\"Optional search query to filter knowledge base data.\",\n tool_mode=True,\n dynamic=True,\n show=False,\n ),\n IntInput(\n name=\"top_k\",\n display_name=\"Top K Results\",\n info=\"Number of top results to return from the knowledge base.\",\n value=5,\n advanced=True,\n required=False,\n dynamic=True,\n show=False,\n ),\n BoolInput(\n name=\"include_metadata\",\n display_name=\"Include Metadata\",\n info=\"Whether to include all metadata in the output. If false, only content is returned.\",\n value=True,\n advanced=False,\n dynamic=True,\n show=False,\n ),\n BoolInput(\n name=\"include_embeddings\",\n display_name=\"Include Embeddings\",\n info=\"Whether to include embeddings in the output. Only applicable if 'Include Metadata' is enabled.\",\n value=False,\n advanced=True,\n dynamic=True,\n show=False,\n ),\n MessageTextInput(\n name=\"metadata_filter\",\n display_name=\"Metadata Filter\",\n info=(\n \"Optional JSON object of user-metadata key/value pairs. Only chunks \"\n 'whose source_metadata matches every key are returned (e.g. {\"tag\": \"invoice\"} '\n 'or {\"tag\": [\"invoice\", \"audit\"]} for OR-of-values). Backends without '\n \"native filtering apply the match client-side after retrieval.\"\n ),\n advanced=True,\n dynamic=True,\n show=False,\n ),\n ]\n\n # ------ Outputs -------------------------------------------------------\n # Both outputs are declared at the class level so the runtime can\n # dispatch to either method depending on which one the saved flow has\n # wired up. ``update_outputs`` filters the canvas-visible output per\n # selected mode; see ``TypeConverterComponent`` and ``MemoryComponent``\n # for the same pattern.\n #\n # Output names match the legacy ``KnowledgeIngestionComponent``\n # (``dataframe_output``) and ``KnowledgeBaseComponent``\n # (``retrieve_data``) so saved flow edges keyed on those names resolve\n # cleanly against the merged component.\n outputs = [\n Output(\n display_name=\"Results\",\n name=\"dataframe_output\",\n method=\"build_kb_info\",\n types=[\"JSON\"],\n selected=\"JSON\",\n ),\n Output(\n display_name=\"Results\",\n name=\"retrieve_data\",\n method=\"retrieve_data\",\n info=\"Returns the data from the selected knowledge base.\",\n types=[\"Table\"],\n selected=\"Table\",\n ),\n ]\n\n # ------ Mode-driven UI updates ---------------------------------------\n async def update_frontend_node(self, new_frontend_node: dict, current_frontend_node: dict):\n \"\"\"Sync the visible output with the current ``mode`` value on canvas load.\n\n Saved flows hit this path: ``update_outputs`` is normally only triggered\n by ``real_time_refresh`` field edits, so without this re-sync the\n canvas could land with both outputs visible.\n \"\"\"\n await super().update_frontend_node(new_frontend_node, current_frontend_node)\n mode_value = new_frontend_node.get(\"template\", {}).get(\"mode\", {}).get(\"value\", MODE_INGEST)\n self.update_outputs(new_frontend_node, \"mode\", mode_value)\n return new_frontend_node\n\n def update_outputs(self, frontend_node: dict, field_name: str, field_value: Any) -> dict:\n \"\"\"Filter visible outputs to match the selected mode.\n\n Triggered by the ``mode`` ``TabInput`` ``real_time_refresh`` flag.\n The class-level ``outputs`` list always carries both entries so the\n runtime can resolve either ``build_kb_info`` or ``retrieve_data``\n regardless of canvas visibility โ€” we just hide the unused one here.\n \"\"\"\n if field_name != \"mode\":\n return frontend_node\n if _is_retrieve_mode(field_value):\n frontend_node[\"outputs\"] = [\n Output(\n display_name=\"Results\",\n name=\"retrieve_data\",\n method=\"retrieve_data\",\n info=\"Returns the data from the selected knowledge base.\",\n types=[\"Table\"],\n selected=\"Table\",\n )\n ]\n else:\n frontend_node[\"outputs\"] = [\n Output(\n display_name=\"Results\",\n name=\"dataframe_output\",\n method=\"build_kb_info\",\n types=[\"JSON\"],\n selected=\"JSON\",\n )\n ]\n return frontend_node\n\n async def update_build_config(\n self,\n build_config,\n field_value: Any,\n field_name: str | None = None,\n ):\n \"\"\"Refresh KB options, drive the create-KB dialog, and hide off-mode fields.\"\"\"\n # Astra-cloud gate covers both ingest and retrieve paths.\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n\n # Always populate the create-KB dialog's embedding-model options so the\n # ModelInput renders correctly regardless of which input triggered the\n # refresh.\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 # KB-picker refresh + create-new-KB flow (lifted from the legacy ingestion\n # component verbatim; relied on by both the canvas refresh button and the\n # dialog-submit path).\n if field_name == \"knowledge_base\":\n # Lazy import keeps lfx importable without langflow installed.\n from langflow.services.database.models.user.crud import get_user_by_id\n\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 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 model_selection = field_value[\"02_embedding_model\"]\n if isinstance(model_selection, dict):\n model_selection = [model_selection]\n\n backend_type, backend_config = self._normalize_backend_selection(\n field_value.get(\"03_knowledge_backend\")\n )\n\n embed_model = get_embeddings(\n model=model_selection,\n user_id=self.user_id,\n )\n\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 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 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 backend_type=backend_type,\n backend_config=backend_config,\n )\n await self._create_knowledge_base_record(\n user_id=self.user_id,\n name=field_value[\"01_new_kb_name\"],\n model_selection=model_selection,\n backend_type=backend_type,\n backend_config=backend_config,\n )\n\n build_config[\"knowledge_base\"][\"options\"] = await get_knowledge_bases(\n _get_knowledge_bases_root_path(),\n user_id=self.user_id,\n )\n if build_config[\"knowledge_base\"][\"value\"] not in build_config[\"knowledge_base\"][\"options\"]:\n build_config[\"knowledge_base\"][\"value\"] = None\n\n # Honor the current mode regardless of which field triggered the refresh.\n # Falls back to MODE_INGEST when ``mode`` is missing (legacy nodes).\n current_mode = build_config.get(\"mode\", {}).get(\"value\") if isinstance(build_config, dict) else None\n if field_name == \"mode\":\n current_mode = field_value\n # Map legacy/emoji-prefixed labels onto the current canonical values so\n # flows saved before the label change still toggle visibility correctly.\n if _is_retrieve_mode(current_mode):\n current_mode = MODE_RETRIEVE\n elif current_mode not in self.mode_config:\n current_mode = MODE_INGEST\n return set_current_fields(\n build_config=build_config if isinstance(build_config, dotdict) else dotdict(build_config),\n action_fields=self.mode_config,\n selected_action=current_mode,\n default_fields=self.default_keys,\n func=set_field_display,\n )\n\n # =====================================================================\n # INGESTION CODE PATH\n # =====================================================================\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\n the result to a single scalar ``bool``.\n \"\"\"\n result = pd.notna(value)\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 config_list = self.column_config if isinstance(self.column_config, list) else []\n\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 backend_type: str = BackendType.CHROMA.value,\n backend_config: dict[str, Any] | None = None,\n ) -> dict[str, Any]:\n \"\"\"Build embedding model metadata from a model selection dict.\"\"\"\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,\n \"api_key\": encrypted_api_key,\n \"api_key_used\": bool(api_key),\n \"chunk_size\": self.chunk_size,\n \"backend_type\": backend_type,\n \"backend_config\": backend_config or {},\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 backend_type: str | None = None,\n backend_config: dict[str, Any] | None = None,\n ) -> None:\n \"\"\"Save embedding model metadata.\"\"\"\n metadata_path = kb_path / \"embedding_metadata.json\"\n existing_metadata: dict[str, Any] = {}\n if metadata_path.exists():\n try:\n existing_metadata = json.loads(metadata_path.read_text())\n except (OSError, json.JSONDecodeError):\n existing_metadata = {}\n\n embedding_metadata = self._build_embedding_metadata(\n model_selection,\n api_key,\n backend_type=backend_type or existing_metadata.get(\"backend_type\") or BackendType.CHROMA.value,\n backend_config=backend_config\n if backend_config is not None\n else existing_metadata.get(\"backend_config\") or {},\n )\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 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 @staticmethod\n def _extract_source_types_from_df(df_source: pd.DataFrame) -> set[str]:\n \"\"\"Pull file extensions out of common path/name columns on the source DataFrame.\n\n The direct-upload ingestion path stores extensions in\n ``embedding_metadata.json[source_types]`` so the KB list can render the\n correct file-type icon. When ingestion happens via a connected\n ``input_df`` (e.g. File โ†’ Knowledge) the same field stayed empty and\n the icon defaulted to a blank tile. We look at the well-known columns\n the File / S3 / cloud-storage components produce and collect any\n plausible extension so the icon is consistent across both flows.\n \"\"\"\n candidate_columns = (\"file_path\", \"file_name\", \"filename\", \"source\", \"path\", \"mimetype\")\n extensions: set[str] = set()\n for col in candidate_columns:\n if col not in df_source.columns:\n continue\n for value in df_source[col].dropna():\n ext = KnowledgeComponent._extension_from_value(value)\n # Drop anything that doesn't look like an extension (URL\n # query strings, version segments, etc.) โ€” the icon palette\n # keys off short alphanumeric tokens like \"pdf\"/\"docx\".\n if ext:\n extensions.add(ext)\n return extensions\n\n @staticmethod\n def _extension_from_value(value: Any) -> str | None:\n \"\"\"Return a normalized file-extension token from a path / filename / MIME string.\n\n Accepts values like ``\"report.PDF\"``, ``\"/docs/notes.txt\"``, or\n ``\"application/pdf\"`` and returns ``\"pdf\"`` / ``\"txt\"``. Returns\n ``None`` if no plausible extension can be derived.\n \"\"\"\n extension_length_limit = 10\n if value is None:\n return None\n text = str(value).strip()\n if not text:\n return None\n # MIME types like ``application/pdf`` carry the canonical extension\n # in the subtype slot โ€” preserve them so File / S3 messages keyed\n # only on ``mimetype`` still resolve to an icon.\n if \"/\" in text and \".\" not in text.rsplit(\"/\", 1)[-1]:\n subtype = text.rsplit(\"/\", 1)[-1].strip().lower()\n return subtype if subtype and len(subtype) <= extension_length_limit and subtype.isalnum() else None\n if \".\" not in text:\n return None\n ext = text.rsplit(\".\", 1)[-1].strip().lower()\n if ext and len(ext) <= extension_length_limit and ext.isalnum():\n return ext\n return None\n\n @classmethod\n def _extract_source_types_from_mapping(cls, mapping: Any) -> set[str]:\n \"\"\"Pull extensions from a Message/Data-style ``data`` dict or a plain dict.\"\"\"\n if not isinstance(mapping, dict):\n return set()\n candidate_keys = (\"file_path\", \"file_name\", \"filename\", \"source\", \"path\", \"mimetype\")\n extensions: set[str] = set()\n for key in candidate_keys:\n ext = cls._extension_from_value(mapping.get(key))\n if ext:\n extensions.add(ext)\n return extensions\n\n @classmethod\n def _extract_source_types_from_input(cls, input_value: Any) -> set[str]:\n \"\"\"Pull file extensions out of the raw component input.\n\n ``convert_to_dataframe`` strips Message/Data fields down to ``text``\n when projecting onto a DataFrame, so file metadata attached to a\n File-component \"Raw Content\" output never reaches\n ``_extract_source_types_from_df``. Looking at the raw input first\n keeps the KB icon consistent with direct upload.\n \"\"\"\n if input_value is None:\n return set()\n if isinstance(input_value, list):\n extensions: set[str] = set()\n for item in input_value:\n extensions |= cls._extract_source_types_from_input(item)\n return extensions\n if isinstance(input_value, pd.DataFrame):\n return cls._extract_source_types_from_df(input_value)\n # Message / Data / JSON all expose a ``data`` dict via the lfx schema.\n mapping = getattr(input_value, \"data\", None)\n if mapping is None and isinstance(input_value, dict):\n mapping = input_value\n return cls._extract_source_types_from_mapping(mapping)\n\n def _merge_source_types(self, kb_path: Path, extensions: set[str]) -> None:\n \"\"\"Merge newly observed extensions into the KB's ``source_types`` metadata.\n\n Mirrors the direct-upload path in ``KBIngestionHelper`` so the icon\n rendering on the Knowledge Bases list works regardless of which\n ingestion route was used.\n \"\"\"\n if not extensions:\n return\n metadata_path = kb_path / \"embedding_metadata.json\"\n if not metadata_path.exists():\n return\n try:\n metadata = json.loads(metadata_path.read_text())\n existing = set(metadata.get(\"source_types\") or [])\n metadata[\"source_types\"] = sorted(existing | extensions)\n metadata_path.write_text(json.dumps(metadata, indent=2))\n except (OSError, ValueError, TypeError, json.JSONDecodeError) as e:\n self.log(f\"Warning: Could not update source_types metadata: {e}\")\n\n async def _update_backend_metadata_metrics(self, kb_path: Path, backend: BaseVectorStoreBackend) -> None:\n \"\"\"Update metadata metrics for non-Chroma backends.\"\"\"\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 chunks = await backend.count()\n characters = 0\n words = 0\n async for batch in backend.iter_documents():\n for document in batch:\n characters += len(document.content)\n words += len(document.content.split())\n\n metadata[\"chunks\"] = chunks\n metadata[\"characters\"] = characters\n metadata[\"words\"] = words\n metadata[\"avg_chunk_size\"] = characters / chunks if chunks else 0.0\n metadata[\"size\"] = await backend.storage_size_bytes()\n metadata_path.write_text(json.dumps(metadata, indent=2))\n except (OSError, ValueError, TypeError, json.JSONDecodeError) as e:\n self.log(f\"Warning: Could not update backend metadata metrics: {e}\")\n\n @staticmethod\n def _normalize_backend_selection(value: Any) -> tuple[str, dict[str, Any]]:\n \"\"\"Normalize a DBProviderInput value into backend type/config.\"\"\"\n if not value:\n return BackendType.CHROMA.value, {}\n\n if isinstance(value, str):\n backend_type = value if value == BackendType.OPENSEARCH.value else BackendType.CHROMA.value\n return (\n backend_type,\n _DEFAULT_OPENSEARCH_CONFIG.copy() if backend_type == BackendType.OPENSEARCH.value else {},\n )\n\n if not isinstance(value, dict):\n return BackendType.CHROMA.value, {}\n\n backend_type = str(value.get(\"backend_type\") or value.get(\"id\") or BackendType.CHROMA.value)\n\n if backend_type == BackendType.OPENSEARCH.value:\n backend_config = value.get(\"backend_config\") or value.get(\"config\") or {}\n if not isinstance(backend_config, dict):\n backend_config = {}\n return BackendType.OPENSEARCH.value, {**_DEFAULT_OPENSEARCH_CONFIG, **backend_config}\n\n if backend_type == \"chroma_cloud\":\n backend_config = value.get(\"backend_config\") or value.get(\"config\") or {}\n if not isinstance(backend_config, dict):\n backend_config = {}\n return BackendType.CHROMA.value, {**_DEFAULT_CHROMA_CLOUD_CONFIG, **backend_config}\n\n return BackendType.CHROMA.value, {}\n\n @staticmethod\n def _get_backend_from_metadata(kb_path: Path) -> tuple[str, dict[str, Any]]:\n metadata_path = kb_path / \"embedding_metadata.json\"\n if not metadata_path.exists():\n return BackendType.CHROMA.value, {}\n try:\n metadata = json.loads(metadata_path.read_text())\n except (OSError, json.JSONDecodeError):\n return BackendType.CHROMA.value, {}\n\n backend_type = str(metadata.get(\"backend_type\") or BackendType.CHROMA.value)\n backend_config = metadata.get(\"backend_config\") or {}\n if not isinstance(backend_config, dict):\n backend_config = {}\n return backend_type, backend_config\n\n async def _create_knowledge_base_record(\n self,\n *,\n user_id: Any,\n name: str,\n model_selection: list[dict[str, Any]],\n backend_type: str,\n backend_config: dict[str, Any],\n ) -> None:\n \"\"\"Persist the component-created KB in the DB when Langflow is available.\"\"\"\n try:\n from langflow.api.utils import knowledge_base_service\n except ImportError:\n return\n\n try:\n await knowledge_base_service.create_record(\n user_id=user_id,\n name=name,\n model_selection=model_selection,\n column_config=self.column_config if isinstance(self.column_config, list) else [],\n backend_type=backend_type,\n backend_config=backend_config,\n )\n except Exception as exc: # noqa: BLE001\n self.log(f\"Warning: could not persist knowledge base record: {exc}\")\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 kb_path.mkdir(parents=True, exist_ok=True)\n\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 metadata[\"columns\"].append(\n {\n \"name\": col_name,\n \"vectorize\": vectorize,\n \"identifier\": identifier,\n }\n )\n\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 ) -> BaseVectorStoreBackend:\n \"\"\"Create vector store using the configured DB provider.\"\"\"\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 backend_type, backend_config = self._get_backend_from_metadata(vector_store_dir)\n backend = create_backend(\n backend_type,\n kb_name=self.knowledge_base,\n kb_path=vector_store_dir,\n backend_config=backend_config,\n embedding_function=embedding_function,\n user_id=self.user_id,\n )\n await backend.ensure_ready()\n\n existing_ids = None\n if backend_type != BackendType.CHROMA.value and not self.allow_duplicates:\n existing_ids = set()\n async for batch in backend.iter_documents():\n for document in batch:\n doc_id = document.metadata.get(\"_id\")\n if doc_id:\n existing_ids.add(doc_id)\n\n data_objects = await self._convert_df_to_data_objects(df_source, config_list, existing_ids=existing_ids)\n\n user_metadata_tag = self._resolve_user_metadata_tag()\n\n documents = []\n for data_obj in data_objects:\n doc = data_obj.to_lc_document()\n if user_metadata_tag:\n doc.metadata[\"source_metadata\"] = user_metadata_tag\n documents.append(doc)\n\n if documents:\n await backend.add_documents(documents)\n self.log(f\"Added {len(documents)} documents to vector store '{self.knowledge_base}'\")\n\n return backend\n\n async def _convert_df_to_data_objects(\n self,\n df_source: pd.DataFrame,\n config_list: list[dict[str, Any]],\n existing_ids: set[str] | None = None,\n ) -> list[Data]:\n \"\"\"Convert DataFrame to Data objects for vector store.\"\"\"\n data_objects: list[Data] = []\n\n if existing_ids is None:\n kb_path = await self._kb_path()\n\n chroma = Chroma(\n persist_directory=str(kb_path),\n collection_name=self.knowledge_base,\n **chroma_langchain_collection_kwargs(),\n )\n\n all_docs = chroma.get()\n\n existing_ids = {metadata.get(\"_id\") for metadata in all_docs[\"metadatas\"] if metadata.get(\"_id\")}\n\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 if identifier:\n identifier_cols.append(col_name)\n\n for _, row in df_source.iterrows():\n identifier_parts = [str(row[col]) for col in content_cols if col in row and self._scalar_notna(row[col])]\n\n page_content = \" \".join(identifier_parts)\n\n data_dict = {\n \"text\": page_content,\n }\n\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 for col in df_source.columns:\n if col not in content_cols and col in row and self._scalar_notna(row[col]):\n value = row[col]\n data_dict[col] = str(value)\n\n page_content_hash = hashlib.sha256(page_content.encode()).hexdigest()\n data_dict[\"_id\"] = page_content_hash\n\n if not self.allow_duplicates and page_content_hash in existing_ids:\n self.log(f\"Skipping duplicate row with hash {page_content_hash}\")\n continue\n\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 \"\"\"Validate collection name.\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 if not (min_length <= len(name) <= max_length):\n return False\n\n if not (name[0].isalnum() and name[-1].isalnum()):\n return False\n\n return re.match(r\"^[a-zA-Z0-9_-]+$\", name) is not None\n\n async def _kb_path(self) -> Path | None:\n cached_path = getattr(self, \"_cached_kb_path\", None)\n if cached_path is not None:\n return cached_path\n\n # Lazy import to keep ``lfx`` importable standalone โ€” langflow's\n # user/DB models are not always available at module load time.\n from langflow.services.database.models.user.crud import get_user_by_id\n\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 self._cached_kb_path = kb_root / kb_user / self.knowledge_base\n\n return self._cached_kb_path\n\n def _resolve_user_metadata_tag(self) -> str:\n \"\"\"Return the JSON-encoded user metadata tag for chunk writes.\"\"\"\n raw = getattr(self, \"metadata_json\", None)\n if not raw:\n return \"\"\n text = raw.strip() if isinstance(raw, str) else raw\n if not text:\n return \"\"\n try:\n decoded = json.loads(text)\n except (TypeError, json.JSONDecodeError) as exc:\n self.log(f\"KnowledgeComponent: metadata_json is not valid JSON ({exc}); skipping metadata stamp.\")\n return \"\"\n if not isinstance(decoded, dict):\n self.log(\"KnowledgeComponent: metadata_json must decode to a JSON object; skipping metadata stamp.\")\n return \"\"\n return json.dumps(decoded, sort_keys=True)\n\n async def build_kb_info(self) -> Data:\n \"\"\"Main ingestion routine โ†’ returns a dict with KB metadata.\n\n The annotation is intentionally narrowed to ``Data`` even though the\n cross-mode fallback below may return a ``DataFrame`` from\n ``retrieve_data``. The frontend builds React-Flow handle IDs from\n this output's type list; widening it to ``Data | DataFrame`` makes\n the API advertise ``[\"JSON\", \"Table\"]`` for the ingest output and\n breaks every saved-edge sourceHandle that was generated against\n a single-type handle (BUG-02). Python doesn't enforce return\n annotations at runtime, so the rare fallback path keeps working.\n \"\"\"\n if _is_retrieve_mode(getattr(self, \"mode\", MODE_INGEST)):\n return await self.retrieve_data()\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n\n run_id: uuid.UUID | None = None\n run_job_id: uuid.UUID | None = None\n run_summary: IngestionSummary | None = None\n run_status: IngestionRunStatus = IngestionRunStatus.SUCCEEDED\n run_error: str | None = None\n kb_record_id: uuid.UUID | None = None\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 config_list = self._validate_column_config(df_source)\n column_metadata = self._build_column_metadata(config_list, df_source)\n\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 if metadata_path.exists():\n settings_service = get_settings_service()\n stored_metadata = json.loads(metadata_path.read_text())\n\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 embedding_model_name = stored_metadata.get(\"embedding_model\")\n embedding_provider = stored_metadata.get(\"embedding_provider\", \"Unknown\")\n if embedding_model_name:\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 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 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 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 run_id, run_job_id, run_summary, kb_record_id = await self._begin_ingestion_run(kb_path)\n if kb_record_id is not None:\n await self._record_kb_status(kb_record_id, \"ingesting\")\n\n backend = await self._create_vector_store(df_source, config_list, embedding_function=embedding_function)\n\n self._save_kb_files(kb_path, config_list)\n\n try:\n if not isinstance(backend, BaseVectorStoreBackend):\n pass\n elif backend.backend_type == BackendType.CHROMA and hasattr(backend, \"raw_langchain_store\"):\n self._update_metadata_metrics(kb_path, backend.raw_langchain_store())\n else:\n await self._update_backend_metadata_metrics(kb_path, backend)\n # Stamp the KB with the file extensions we just ingested so\n # the Knowledge Bases list renders the correct icon for\n # flow-driven ingestion (input_df), matching direct upload.\n # We look at the raw input first because ``convert_to_dataframe``\n # drops Message/Data metadata fields (file_path, mimetype, โ€ฆ)\n # when projecting onto the DataFrame.\n source_types = self._extract_source_types_from_input(input_value)\n source_types |= self._extract_source_types_from_df(df_source)\n self._merge_source_types(kb_path, source_types)\n finally:\n if isinstance(backend, BaseVectorStoreBackend):\n await backend.teardown()\n\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 if run_summary is not None:\n run_summary.record_item(\n IngestionItemResult(\n item_id=self.knowledge_base,\n display_name=f\"{self.knowledge_base} ({len(df_source)} rows)\",\n status=IngestionItemStatus.SUCCEEDED,\n chunks_created=len(df_source),\n ),\n )\n\n if kb_record_id is not None:\n await self._record_kb_stats(kb_record_id, kb_path)\n await self._record_kb_status(kb_record_id, \"ready\")\n\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 run_status = IngestionRunStatus.FAILED\n run_error = str(e) or e.__class__.__name__\n msg = f\"Error during KB ingestion: {e}\"\n raise RuntimeError(msg) from e\n except Exception as e:\n run_status = IngestionRunStatus.FAILED\n run_error = str(e) or e.__class__.__name__\n raise\n finally:\n if run_id is not None and run_summary is not None:\n await self._finalize_ingestion_run(\n run_id=run_id,\n job_id=run_job_id,\n summary=run_summary,\n status=run_status,\n error_message=run_error,\n )\n if kb_record_id is not None and run_status is IngestionRunStatus.FAILED:\n await self._record_kb_status(kb_record_id, \"failed\", failure_reason=run_error)\n\n async def _begin_ingestion_run(\n self,\n kb_path: Path,\n ) -> tuple[uuid.UUID | None, uuid.UUID | None, IngestionSummary | None, uuid.UUID | None]:\n \"\"\"Create a parent ``Job`` and seed an ingestion-run row.\"\"\"\n if not self.user_id:\n self.log(\"No user_id on component; skipping ingestion-run tracking.\")\n return None, None, None, None\n\n try:\n user_uuid = uuid.UUID(str(self.user_id))\n except (ValueError, TypeError) as exc:\n self.log(f\"Could not coerce user_id={self.user_id!r} to UUID; skipping run tracking ({exc}).\")\n return None, None, None, None\n\n try:\n from langflow.api.utils import ingestion_run_service, knowledge_base_service\n from langflow.services.database.models.jobs.model import JobStatus, JobType\n from langflow.services.deps import get_job_service\n except ImportError as exc:\n self.log(f\"Run-history wiring unavailable; ingestion will not be recorded ({exc}).\")\n return None, None, None, None\n\n try:\n kb_record = await knowledge_base_service.get_by_user_and_name(user_uuid, self.knowledge_base)\n kb_record_id = kb_record.id if kb_record is not None else None\n\n user_metadata = self._parse_user_metadata_dict()\n\n job_id = uuid.uuid4()\n job_service = get_job_service()\n raw_flow_id = getattr(self, \"flow_id\", None)\n flow_id_uuid = uuid.UUID(str(raw_flow_id)) if raw_flow_id else job_id\n await job_service.create_job(\n job_id=job_id,\n flow_id=flow_id_uuid,\n job_type=JobType.INGESTION,\n asset_id=kb_record_id,\n asset_type=\"knowledge_base\",\n user_id=user_uuid,\n )\n await job_service.update_job_status(job_id, JobStatus.IN_PROGRESS)\n\n source = FlowComponentSource(\n user_id=user_uuid,\n source_config={\n \"knowledge_base\": self.knowledge_base,\n \"kb_path\": str(kb_path),\n \"flow_id\": str(flow_id_uuid),\n },\n )\n\n run_id = await ingestion_run_service.create_run(\n kb_name=self.knowledge_base,\n source=source,\n job_id=job_id,\n user_id=user_uuid,\n kb_id=kb_record_id,\n user_metadata=user_metadata,\n )\n if run_id is None:\n return None, job_id, None, kb_record_id\n await ingestion_run_service.mark_running(run_id)\n\n summary = IngestionSummary(\n kb_name=self.knowledge_base,\n source_type=source.source_type.value,\n user_id=user_uuid,\n job_id=job_id,\n source_config=source.describe().get(\"config\") or {},\n user_metadata=user_metadata,\n )\n self.log(f\"Started ingestion run job_id={job_id} kb_name={self.knowledge_base} kb_id={kb_record_id}\")\n except Exception as exc: # noqa: BLE001 โ€” telemetry must never abort ingestion\n self.log(f\"Could not begin ingestion-run tracking: {exc}\")\n return None, None, None, None\n else:\n return run_id, job_id, summary, kb_record_id\n\n async def _finalize_ingestion_run(\n self,\n *,\n run_id: uuid.UUID,\n job_id: uuid.UUID | None,\n summary: IngestionSummary,\n status: IngestionRunStatus,\n error_message: str | None,\n ) -> None:\n \"\"\"Persist the final summary and transition the parent Job.\"\"\"\n try:\n from langflow.api.utils import ingestion_run_service\n from langflow.services.database.models.jobs.model import JobStatus\n from langflow.services.deps import get_job_service\n except ImportError as exc:\n self.log(f\"Run-history wiring unavailable; ingestion-run finalize skipped ({exc}).\")\n return\n\n try:\n await ingestion_run_service.finalize_run(\n run_id,\n summary=summary,\n status=status,\n error_message=error_message,\n )\n if job_id is not None:\n terminal_status = JobStatus.COMPLETED if status is not IngestionRunStatus.FAILED else JobStatus.FAILED\n await get_job_service().update_job_status(job_id, terminal_status, finished_timestamp=True)\n except Exception as exc: # noqa: BLE001 โ€” telemetry must never re-raise\n self.log(f\"Could not finalize ingestion-run tracking: {exc}\")\n\n async def _record_kb_status(\n self,\n kb_record_id: uuid.UUID,\n status_value: str,\n *,\n failure_reason: str | None = None,\n ) -> None:\n \"\"\"Mirror Path A's KB-row status transitions.\"\"\"\n try:\n from langflow.api.utils import knowledge_base_service\n from langflow.services.database.models.knowledge_base.model import KnowledgeBaseStatus\n except ImportError:\n return\n try:\n await knowledge_base_service.update_status(\n kb_record_id,\n status=KnowledgeBaseStatus(status_value),\n failure_reason=failure_reason,\n )\n except Exception as exc: # noqa: BLE001\n self.log(f\"Could not update KB status to {status_value}: {exc}\")\n\n async def _record_kb_stats(self, kb_record_id: uuid.UUID, kb_path: Path) -> None:\n \"\"\"Push freshly-refreshed metrics from embedding_metadata.json onto the DB row.\"\"\"\n try:\n from langflow.api.utils import knowledge_base_service\n except ImportError:\n self.log(\"knowledge_base_service unavailable; KB stats will not sync to DB row.\")\n return\n metadata_path = kb_path / \"embedding_metadata.json\"\n if not metadata_path.exists():\n self.log(f\"No embedding_metadata.json at {metadata_path}; skipping KB stats sync.\")\n return\n try:\n metadata = json.loads(metadata_path.read_text())\n except (OSError, json.JSONDecodeError) as exc:\n self.log(f\"Could not read KB metadata for stats sync: {exc}\")\n return\n chunks = int(metadata.get(\"chunks\", 0) or 0)\n words = int(metadata.get(\"words\", 0) or 0)\n characters = int(metadata.get(\"characters\", 0) or 0)\n try:\n await knowledge_base_service.update_stats(\n kb_record_id,\n chunks=chunks,\n words=words,\n characters=characters,\n size_bytes=int(metadata.get(\"size\", 0) or 0),\n source_types=list(metadata.get(\"source_types\") or []),\n chunk_size=metadata.get(\"chunk_size\"),\n chunk_overlap=metadata.get(\"chunk_overlap\"),\n separator=metadata.get(\"separator\"),\n )\n self.log(f\"Synced KB stats to DB row {kb_record_id}: chunks={chunks} words={words} characters={characters}\")\n except Exception as exc: # noqa: BLE001\n self.log(f\"Could not sync KB stats to DB row {kb_record_id}: {exc}\")\n\n def _parse_user_metadata_dict(self) -> dict[str, Any]:\n \"\"\"Decode ``metadata_json`` to a dict, or ``{}`` on any error.\"\"\"\n raw = getattr(self, \"metadata_json\", None)\n if not raw:\n return {}\n text = raw.strip() if isinstance(raw, str) else raw\n if not text:\n return {}\n try:\n decoded = json.loads(text)\n except (TypeError, json.JSONDecodeError):\n return {}\n return decoded if isinstance(decoded, dict) else {}\n\n # =====================================================================\n # RETRIEVAL CODE PATH\n # =====================================================================\n @property\n def _user_uuid(self) -> uuid.UUID | None:\n \"\"\"Return self.user_id as a UUID, converting from str if necessary.\"\"\"\n if not self.user_id:\n return None\n return self.user_id if isinstance(self.user_id, uuid.UUID) else uuid.UUID(self.user_id)\n\n def _get_kb_metadata(self, kb_path: Path) -> dict:\n \"\"\"Load the knowledge base's embedding metadata file.\"\"\"\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n return load_kb_metadata(kb_path, log_label=f\"knowledge base '{self.knowledge_base}'\")\n\n async def _resolve_backend(self, *, kb_user: str) -> tuple[str, dict[str, Any]]: # noqa: ARG002\n \"\"\"Return ``(backend_type, backend_config)`` for this KB.\"\"\"\n try:\n from langflow.api.utils import knowledge_base_service\n\n user_uuid = self._user_uuid\n if user_uuid is None:\n return BackendType.CHROMA.value, {}\n record = await knowledge_base_service.get_by_user_and_name(user_uuid, self.knowledge_base)\n except Exception as exc: # noqa: BLE001\n logger.debug(\"KB record lookup failed: %s\", exc)\n return BackendType.CHROMA.value, {}\n\n if record is None:\n return BackendType.CHROMA.value, {}\n return (\n record.backend_type or BackendType.CHROMA.value,\n record.backend_config or {},\n )\n\n def _resolve_model_selection(self, metadata: dict[str, Any]) -> list[dict[str, Any]]:\n \"\"\"Resolve the ``get_embeddings``-compatible model selection from metadata.\"\"\"\n model_selection = metadata.get(\"model_selection\")\n if model_selection:\n selection_list = [model_selection] if isinstance(model_selection, dict) else list(model_selection)\n return [self._hydrate_model_metadata(entry) for entry in selection_list]\n\n embedding_model_name = metadata.get(\"embedding_model\")\n embedding_provider = metadata.get(\"embedding_provider\", \"Unknown\")\n if not embedding_model_name:\n msg = (\n f\"Knowledge base '{self.knowledge_base}' has no embedding model recorded; \"\n \"re-create it with a supported embedding model.\"\n )\n raise ValueError(msg)\n\n match = self._find_catalog_entry(embedding_model_name)\n if match is None:\n msg = (\n f\"Embedding model '{embedding_model_name}' (provider '{embedding_provider}') \"\n \"recorded for this knowledge base is no longer available in the model registry. \"\n \"Please re-create the knowledge base with a supported embedding model.\"\n )\n raise ValueError(msg)\n return [match]\n\n def _hydrate_model_metadata(self, entry: dict[str, Any]) -> dict[str, Any]:\n \"\"\"Fill in ``metadata.embedding_class`` / ``param_mapping`` if missing.\"\"\"\n entry_metadata = entry.get(\"metadata\") or {}\n has_class = bool(entry_metadata.get(\"embedding_class\"))\n has_mapping = bool(entry_metadata.get(\"param_mapping\"))\n if has_class and has_mapping:\n return entry\n\n model_name = entry.get(\"name\")\n if not model_name:\n return entry\n\n catalog_entry = self._find_catalog_entry(model_name)\n if catalog_entry is None:\n return entry\n\n catalog_metadata = catalog_entry.get(\"metadata\") or {}\n merged_metadata = {**catalog_metadata, **entry_metadata}\n return {**entry, \"metadata\": merged_metadata}\n\n def _find_catalog_entry(self, model_name: str) -> dict[str, Any] | None:\n \"\"\"Look up an embedding model by name in the unified-models catalog.\"\"\"\n options = get_embedding_model_options(user_id=self.user_id)\n return next((o for o in options if o.get(\"name\") == model_name), None)\n\n async def retrieve_data(self) -> DataFrame:\n \"\"\"Retrieve data from the selected knowledge base.\n\n Annotation narrowed to ``DataFrame`` to keep this output's handle\n type list at ``[\"Table\"]`` only; widening to a union surfaces\n ``[\"Table\", \"JSON\"]`` on the API and breaks every starter-project\n edge whose sourceHandle was stored against a single-type handle\n (BUG-02). The cross-mode defensive fallback may still return a\n ``Data`` at runtime โ€” Python ignores return annotations, so\n nothing breaks.\n \"\"\"\n if not _is_retrieve_mode(getattr(self, \"mode\", MODE_INGEST)):\n return await self.build_kb_info()\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n\n # Lazy import: langflow's user/DB models aren't part of lfx's\n # standalone install, so ``lfx run .json`` can't resolve\n # this symbol at module import time. Deferring to use keeps the\n # component importable in both environments.\n from langflow.services.database.models.user.crud import get_user_by_id\n\n async with session_scope() as db:\n if not self.user_id:\n msg = \"User ID is required for fetching Knowledge Base data.\"\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 kb_path = _get_knowledge_bases_root_path() / kb_user / self.knowledge_base\n\n metadata = self._get_kb_metadata(kb_path)\n if not metadata:\n msg = f\"Metadata not found for knowledge base: {self.knowledge_base}. Ensure it has been indexed.\"\n raise ValueError(msg)\n\n model_selection = self._resolve_model_selection(metadata)\n chunk_size = metadata.get(\"chunk_size\")\n embedding_function = get_embeddings(\n model=model_selection,\n user_id=self.user_id,\n chunk_size=chunk_size,\n )\n\n backend_type, backend_config = await self._resolve_backend(kb_user=kb_user)\n backend = create_backend(\n backend_type,\n kb_name=self.knowledge_base,\n kb_path=kb_path,\n backend_config=backend_config,\n embedding_function=embedding_function,\n user_id=self.user_id,\n )\n try:\n user_metadata_filter = _parse_metadata_filter(getattr(self, \"metadata_filter\", None))\n use_scores = bool(self.search_query)\n search_k = self.top_k * 4 if user_metadata_filter else self.top_k\n results = await backend.similarity_search(\n query=self.search_query or \"\",\n k=search_k,\n with_scores=use_scores,\n )\n if user_metadata_filter:\n results = [\n (doc, score) for doc, score in results if _chunk_matches_filter(doc.metadata, user_metadata_filter)\n ]\n results = results[: self.top_k]\n\n embeddings_by_key: dict[tuple[str, str], list[float]] = {}\n if self.include_embeddings and results:\n # Join each retrieved chunk to its stored embedding. We key on\n # ``_id`` when present and fall back to page content otherwise,\n # so KBs populated by direct file upload โ€” whose chunks carry no\n # ``_id`` โ€” still resolve. See ``_embedding_match_key``.\n wanted_keys = {_embedding_match_key(doc.page_content, doc.metadata) for doc, _score in results}\n async for batch in backend.iter_documents(include_embeddings=True):\n for entry in batch:\n if entry.embedding is None:\n continue\n key = _embedding_match_key(entry.content, entry.metadata)\n if key in wanted_keys:\n embeddings_by_key[key] = entry.embedding\n if len(embeddings_by_key) == len(wanted_keys):\n break\n\n data_list: list[Data] = []\n for doc, score in results:\n kwargs: dict[str, Any] = {\"content\": doc.page_content}\n if use_scores:\n kwargs[\"_score\"] = -1 * score\n if self.include_metadata:\n kwargs.update(doc.metadata)\n if self.include_embeddings:\n kwargs[\"_embeddings\"] = embeddings_by_key.get(_embedding_match_key(doc.page_content, doc.metadata))\n data_list.append(Data(**kwargs))\n\n return DataFrame(data=data_list)\n finally:\n await backend.teardown()\n\n\ndef _embedding_match_key(content: str, metadata: dict[str, Any] | None) -> tuple[str, str]:\n \"\"\"Build a stable key for aligning a retrieved chunk with its stored embedding.\n\n Embeddings are gathered separately (via ``iter_documents``) and then joined\n back onto the search results. Component-driven ingestion stamps a\n content-hash ``_id`` on every chunk, but direct file-upload ingestion\n (``KBIngestionHelper.perform_ingestion``) does not. Keying the join purely on\n ``_id`` therefore left every upload-populated KB with ``_embeddings: None``.\n\n We prefer ``_id`` when present (so legitimately distinct chunks that happen to\n share text aren't collapsed) and fall back to the chunk's page content\n otherwise. The content fallback is exact for the embedding use case: two\n chunks with identical text necessarily share the same embedding vector. The\n leading ``\"id\"`` / ``\"content\"`` tag namespaces the two key spaces so a mixed\n KB (some chunks with ``_id``, some without) never cross-matches.\n \"\"\"\n doc_id = metadata.get(\"_id\") if metadata else None\n if doc_id:\n return (\"id\", str(doc_id))\n return (\"content\", content or \"\")\n\n\ndef _parse_metadata_filter(raw: str | None) -> dict[str, list[str]]:\n \"\"\"Decode the ``metadata_filter`` input into a {key: [values]} map.\n\n Empty or malformed input maps to an empty filter so retrieval falls back\n to the unfiltered path. JSON errors are swallowed here rather than raised:\n surfacing component-config errors at the canvas node would break a flow\n run for what is meant to be an optional refinement.\n \"\"\"\n if not raw:\n return {}\n text = raw.strip() if isinstance(raw, str) else raw\n if not text:\n return {}\n try:\n decoded = json.loads(text)\n except (TypeError, json.JSONDecodeError):\n logger.warning(\"KnowledgeComponent: metadata_filter is not valid JSON; ignoring filter.\")\n return {}\n if not isinstance(decoded, dict):\n logger.warning(\"KnowledgeComponent: metadata_filter must be a JSON object; ignoring filter.\")\n return {}\n result: dict[str, list[str]] = {}\n for key, value in decoded.items():\n if not isinstance(key, str):\n continue\n if isinstance(value, list):\n result[key] = [str(entry) for entry in value]\n else:\n result[key] = [str(value)]\n return result\n\n\ndef _chunk_matches_filter(metadata: dict[str, Any] | None, filt: dict[str, list[str]]) -> bool:\n \"\"\"AND across keys, OR within key values, mirroring the chunks endpoint.\"\"\"\n if not filt:\n return True\n if not metadata:\n return False\n raw = metadata.get(\"source_metadata\")\n if not raw:\n return False\n try:\n stored = json.loads(raw) if isinstance(raw, str) else raw\n except json.JSONDecodeError:\n return False\n if not isinstance(stored, dict):\n return False\n for key, expected_values in filt.items():\n actual = stored.get(key)\n if actual is None:\n return False\n actual_set = {str(entry) for entry in actual} if isinstance(actual, list) else {str(actual)}\n if not actual_set & set(expected_values):\n return False\n return True\n" + "value": "\"\"\"Unified Knowledge component โ€” ingest into or retrieve from a knowledge base.\n\nThis component merges what used to live in ``ingestion.py`` and\n``retrieval.py`` into a single, mode-driven component. A ``TabInput``\n(\"๐Ÿ“ฅ Ingest\" / \"๐Ÿ” Retrieve\") drives which inputs and which output are\nvisible. The merged shape gives users one node per KB instead of two\nparallel ones that always shared the same ``knowledge_base`` picker,\nembedding-model metadata, and backend registry.\n\nBoth legacy classes (``KnowledgeIngestionComponent``,\n``KnowledgeBaseComponent``) remain importable as thin subclasses of\n``KnowledgeComponent`` โ€” see the sibling ``ingestion.py`` / ``retrieval.py``\nmodules โ€” so saved flows continue to load unchanged.\n\"\"\"\n\nfrom __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 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\n\nfrom lfx.base.knowledge_bases.backends import BackendType, BaseVectorStoreBackend, create_backend\nfrom lfx.base.knowledge_bases.ingestion_sources.base import (\n IngestionItemResult,\n IngestionItemStatus,\n IngestionRunStatus,\n IngestionSummary,\n)\nfrom lfx.base.knowledge_bases.ingestion_sources.flow_component import FlowComponentSource\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.base.vectorstores.chroma_security import chroma_langchain_collection_kwargs\nfrom lfx.components.files_and_knowledge._kb_paths import (\n KBKeyDecryptError,\n load_kb_metadata,\n)\nfrom lfx.components.files_and_knowledge._kb_paths import (\n get_knowledge_bases_root_path as _get_knowledge_bases_root_path,\n)\nfrom lfx.components.processing.converter import convert_to_dataframe\nfrom lfx.custom import Component\nfrom lfx.io import (\n BoolInput,\n DBProviderInput,\n DropdownInput,\n HandleInput,\n IntInput,\n MessageTextInput,\n ModelInput,\n Output,\n SecretStrInput,\n StrInput,\n TabInput,\n TableInput,\n)\nfrom lfx.log.logger import logger\nfrom lfx.schema.data import Data\nfrom lfx.schema.dataframe import DataFrame\nfrom lfx.schema.dotdict import dotdict\nfrom lfx.schema.table import EditMode\nfrom lfx.services.deps import (\n get_settings_service,\n session_scope,\n)\nfrom lfx.utils.component_utils import set_current_fields, set_field_display\nfrom lfx.utils.validate_cloud import raise_error_if_astra_cloud_disable_component\n\n\ndef _inputs_for_mode(default_mode: str) -> list:\n \"\"\"Return a fresh copy of the canonical inputs list with show flags set for the given mode.\n\n Used by the legacy subclasses so a saved flow keyed on\n ``KnowledgeIngestion`` / ``KnowledgeBase`` lands on a node template\n whose default visibility already matches the pinned mode โ€” no\n \"wait for the user to click the mode tab\" UX gap on load.\n \"\"\"\n always_visible = set(KnowledgeComponent.default_keys) | set(KnowledgeComponent.mode_config[default_mode])\n inputs_copy = []\n for inp in KnowledgeComponent.inputs:\n clone = inp.model_copy(deep=True) if hasattr(inp, \"model_copy\") else inp\n if clone.name == \"mode\":\n clone.value = default_mode\n else:\n clone.show = clone.name in always_visible\n inputs_copy.append(clone)\n return inputs_copy\n\n\nif TYPE_CHECKING:\n from pathlib import Path\n\n# Mode constants. Plain-text labels (no emoji) for consistency with the rest of\n# Langflow's TabInput palette โ€” see ``MemoryComponent`` for the same convention.\nMODE_INGEST = \"Ingest\"\nMODE_RETRIEVE = \"Retrieve\"\n\n\ndef _is_retrieve_mode(value: Any) -> bool:\n \"\"\"Lenient mode check: treats any label containing 'Retrieve' as retrieve mode.\n\n Older saved flows may carry the emoji-prefixed labels (\"๐Ÿ“ฅ Ingest\" /\n \"๐Ÿ” Retrieve\") this component used to ship with; substring matching\n keeps those loading without forcing a flow rewrite.\n \"\"\"\n return isinstance(value, str) and \"Retrieve\" in value\n\n\n# Error message used by both the ingest and retrieve paths when the user is\n# running against an Astra cloud environment that disables these flows.\nastra_error_msg = \"Knowledge ingestion and retrieval are not supported in Astra cloud environment.\"\n\n_DEFAULT_OPENSEARCH_CONFIG = {\n \"url_variable\": \"OPENSEARCH_URL\",\n \"username_variable\": \"OPENSEARCH_USERNAME\",\n \"password_variable\": \"OPENSEARCH_PASSWORD\", # pragma: allowlist secret\n \"index_name\": \"\",\n \"vector_field\": \"vector_field\",\n \"text_field\": \"text\",\n}\n\n_DEFAULT_CHROMA_CLOUD_CONFIG = {\n \"mode\": \"cloud\",\n \"tenant_variable\": \"CHROMA_TENANT\",\n \"database_variable\": \"CHROMA_DATABASE\",\n \"api_key_variable\": \"CHROMA_API_KEY\", # pragma: allowlist secret\n}\n\n\nclass KnowledgeComponent(Component):\n \"\"\"One component for both writing into and reading from a Langflow knowledge base.\n\n A ``TabInput`` switches between ingestion and retrieval. The\n ``update_build_config`` / ``update_outputs`` hooks hide the inputs\n and the output that don't apply to the current mode so the canvas\n node stays focused.\n \"\"\"\n\n display_name = \"Knowledge\"\n description = \"Ingest into or retrieve from a Langflow knowledge base.\"\n icon = \"database\"\n name = \"Knowledge\"\n\n # ------ Mode โ†’ visible-fields wiring ---------------------------------\n # ``default_keys`` are inputs always visible regardless of mode.\n # ``mode_config`` lists the inputs unique to each mode; everything outside\n # this set is hidden when its mode is not selected.\n default_keys: list[str] = [\"mode\", \"knowledge_base\"]\n mode_config: dict[str, list[str]] = {\n MODE_INGEST: [\n \"input_df\",\n \"column_config\",\n \"chunk_size\",\n \"api_key\",\n \"allow_duplicates\",\n \"metadata_json\",\n ],\n MODE_RETRIEVE: [\n \"search_query\",\n \"top_k\",\n \"include_metadata\",\n \"include_embeddings\",\n \"metadata_filter\",\n ],\n }\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 Base\",\n \"field_order\": [\n \"01_new_kb_name\",\n \"02_embedding_model\",\n \"03_knowledge_backend\",\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 \"03_knowledge_backend\": DBProviderInput(\n name=\"knowledge_backend\",\n display_name=\"DB Provider\",\n info=(\n \"Select where this knowledge base stores vectors. \"\n \"OpenSearch uses the global DB Providers settings.\"\n ),\n required=True,\n ),\n },\n },\n }\n }\n )\n\n # ------ Inputs --------------------------------------------------------\n # ``knowledge_base`` is kept at position 0 so legacy tests that check\n # ``component.inputs[0].dialog_inputs`` continue to work.\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 TabInput(\n name=\"mode\",\n display_name=\"Mode\",\n options=[MODE_INGEST, MODE_RETRIEVE],\n value=MODE_INGEST,\n info=\"Switch between writing new data into the knowledge base and querying it.\",\n real_time_refresh=True,\n tool_mode=True,\n ),\n # --- Ingest-only inputs (default-shown; hidden when mode == Retrieve) -\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 dynamic=True,\n show=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 dynamic=True,\n show=True,\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 dynamic=True,\n show=True,\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 dynamic=True,\n show=True,\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 dynamic=True,\n show=True,\n ),\n StrInput(\n name=\"metadata_json\",\n display_name=\"Metadata\",\n info=(\n \"Optional JSON object of user metadata applied to every chunk produced by this \"\n 'run (e.g. {\"tag\": \"invoice\", \"year\": \"2026\"}). Same shape as the upload modal '\n \"Metadata section so chunks browser filters + Knowledge retrieval metadata_filter \"\n \"work uniformly across upload, folder, and flow-driven ingestion. Malformed JSON is \"\n \"ignored with a warning rather than failing the run.\"\n ),\n advanced=True,\n required=False,\n dynamic=True,\n show=True,\n ),\n # --- Retrieve-only inputs (default-hidden; shown when mode == Retrieve) -\n MessageTextInput(\n name=\"search_query\",\n display_name=\"Search Query\",\n info=\"Optional search query to filter knowledge base data.\",\n tool_mode=True,\n dynamic=True,\n show=False,\n ),\n IntInput(\n name=\"top_k\",\n display_name=\"Top K Results\",\n info=\"Number of top results to return from the knowledge base.\",\n value=5,\n advanced=True,\n required=False,\n dynamic=True,\n show=False,\n ),\n BoolInput(\n name=\"include_metadata\",\n display_name=\"Include Metadata\",\n info=\"Whether to include all metadata in the output. If false, only content is returned.\",\n value=True,\n advanced=False,\n dynamic=True,\n show=False,\n ),\n BoolInput(\n name=\"include_embeddings\",\n display_name=\"Include Embeddings\",\n info=\"Whether to include embeddings in the output. Only applicable if 'Include Metadata' is enabled.\",\n value=False,\n advanced=True,\n dynamic=True,\n show=False,\n ),\n MessageTextInput(\n name=\"metadata_filter\",\n display_name=\"Metadata Filter\",\n info=(\n \"Optional JSON object of user-metadata key/value pairs. Only chunks \"\n 'whose source_metadata matches every key are returned (e.g. {\"tag\": \"invoice\"} '\n 'or {\"tag\": [\"invoice\", \"audit\"]} for OR-of-values). Backends without '\n \"native filtering apply the match client-side after retrieval.\"\n ),\n advanced=True,\n dynamic=True,\n show=False,\n ),\n ]\n\n # ------ Outputs -------------------------------------------------------\n # Both outputs are declared at the class level so the runtime can\n # dispatch to either method depending on which one the saved flow has\n # wired up. ``update_outputs`` filters the canvas-visible output per\n # selected mode; see ``TypeConverterComponent`` and ``MemoryComponent``\n # for the same pattern.\n #\n # Output names match the legacy ``KnowledgeIngestionComponent``\n # (``dataframe_output``) and ``KnowledgeBaseComponent``\n # (``retrieve_data``) so saved flow edges keyed on those names resolve\n # cleanly against the merged component.\n outputs = [\n Output(\n display_name=\"Results\",\n name=\"dataframe_output\",\n method=\"build_kb_info\",\n types=[\"JSON\"],\n selected=\"JSON\",\n ),\n Output(\n display_name=\"Results\",\n name=\"retrieve_data\",\n method=\"retrieve_data\",\n info=\"Returns the data from the selected knowledge base.\",\n types=[\"Table\"],\n selected=\"Table\",\n ),\n ]\n\n # ------ Mode-driven UI updates ---------------------------------------\n async def update_frontend_node(self, new_frontend_node: dict, current_frontend_node: dict):\n \"\"\"Sync the visible output with the current ``mode`` value on canvas load.\n\n Saved flows hit this path: ``update_outputs`` is normally only triggered\n by ``real_time_refresh`` field edits, so without this re-sync the\n canvas could land with both outputs visible.\n \"\"\"\n await super().update_frontend_node(new_frontend_node, current_frontend_node)\n mode_value = new_frontend_node.get(\"template\", {}).get(\"mode\", {}).get(\"value\", MODE_INGEST)\n self.update_outputs(new_frontend_node, \"mode\", mode_value)\n return new_frontend_node\n\n def update_outputs(self, frontend_node: dict, field_name: str, field_value: Any) -> dict:\n \"\"\"Filter visible outputs to match the selected mode.\n\n Triggered by the ``mode`` ``TabInput`` ``real_time_refresh`` flag.\n The class-level ``outputs`` list always carries both entries so the\n runtime can resolve either ``build_kb_info`` or ``retrieve_data``\n regardless of canvas visibility โ€” we just hide the unused one here.\n \"\"\"\n if field_name != \"mode\":\n return frontend_node\n if _is_retrieve_mode(field_value):\n frontend_node[\"outputs\"] = [\n Output(\n display_name=\"Results\",\n name=\"retrieve_data\",\n method=\"retrieve_data\",\n info=\"Returns the data from the selected knowledge base.\",\n types=[\"Table\"],\n selected=\"Table\",\n )\n ]\n else:\n frontend_node[\"outputs\"] = [\n Output(\n display_name=\"Results\",\n name=\"dataframe_output\",\n method=\"build_kb_info\",\n types=[\"JSON\"],\n selected=\"JSON\",\n )\n ]\n return frontend_node\n\n async def update_build_config(\n self,\n build_config,\n field_value: Any,\n field_name: str | None = None,\n ):\n \"\"\"Refresh KB options, drive the create-KB dialog, and hide off-mode fields.\"\"\"\n # Astra-cloud gate covers both ingest and retrieve paths.\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n\n # Always populate the create-KB dialog's embedding-model options so the\n # ModelInput renders correctly regardless of which input triggered the\n # refresh.\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 # KB-picker refresh + create-new-KB flow (lifted from the legacy ingestion\n # component verbatim; relied on by both the canvas refresh button and the\n # dialog-submit path).\n if field_name == \"knowledge_base\":\n # Lazy import keeps lfx importable without langflow installed.\n from langflow.services.database.models.user.crud import get_user_by_id\n\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 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 model_selection = field_value[\"02_embedding_model\"]\n if isinstance(model_selection, dict):\n model_selection = [model_selection]\n\n backend_type, backend_config = self._normalize_backend_selection(\n field_value.get(\"03_knowledge_backend\")\n )\n\n embed_model = get_embeddings(\n model=model_selection,\n user_id=self.user_id,\n )\n\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 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 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 backend_type=backend_type,\n backend_config=backend_config,\n )\n await self._create_knowledge_base_record(\n user_id=self.user_id,\n name=field_value[\"01_new_kb_name\"],\n model_selection=model_selection,\n backend_type=backend_type,\n backend_config=backend_config,\n )\n\n build_config[\"knowledge_base\"][\"options\"] = await get_knowledge_bases(\n _get_knowledge_bases_root_path(),\n user_id=self.user_id,\n )\n if build_config[\"knowledge_base\"][\"value\"] not in build_config[\"knowledge_base\"][\"options\"]:\n build_config[\"knowledge_base\"][\"value\"] = None\n\n # Honor the current mode regardless of which field triggered the refresh.\n # Falls back to MODE_INGEST when ``mode`` is missing (legacy nodes).\n current_mode = build_config.get(\"mode\", {}).get(\"value\") if isinstance(build_config, dict) else None\n if field_name == \"mode\":\n current_mode = field_value\n # Map legacy/emoji-prefixed labels onto the current canonical values so\n # flows saved before the label change still toggle visibility correctly.\n if _is_retrieve_mode(current_mode):\n current_mode = MODE_RETRIEVE\n elif current_mode not in self.mode_config:\n current_mode = MODE_INGEST\n return set_current_fields(\n build_config=build_config if isinstance(build_config, dotdict) else dotdict(build_config),\n action_fields=self.mode_config,\n selected_action=current_mode,\n default_fields=self.default_keys,\n func=set_field_display,\n )\n\n # =====================================================================\n # INGESTION CODE PATH\n # =====================================================================\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\n the result to a single scalar ``bool``.\n \"\"\"\n result = pd.notna(value)\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 config_list = self.column_config if isinstance(self.column_config, list) else []\n\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 backend_type: str = BackendType.CHROMA.value,\n backend_config: dict[str, Any] | None = None,\n ) -> dict[str, Any]:\n \"\"\"Build embedding model metadata from a model selection dict.\"\"\"\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,\n \"api_key\": encrypted_api_key,\n \"api_key_used\": bool(api_key),\n \"chunk_size\": self.chunk_size,\n \"backend_type\": backend_type,\n \"backend_config\": backend_config or {},\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 backend_type: str | None = None,\n backend_config: dict[str, Any] | None = None,\n ) -> None:\n \"\"\"Save embedding model metadata.\"\"\"\n metadata_path = kb_path / \"embedding_metadata.json\"\n existing_metadata: dict[str, Any] = {}\n if metadata_path.exists():\n try:\n existing_metadata = json.loads(metadata_path.read_text())\n except (OSError, json.JSONDecodeError):\n existing_metadata = {}\n\n embedding_metadata = self._build_embedding_metadata(\n model_selection,\n api_key,\n backend_type=backend_type or existing_metadata.get(\"backend_type\") or BackendType.CHROMA.value,\n backend_config=backend_config\n if backend_config is not None\n else existing_metadata.get(\"backend_config\") or {},\n )\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 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 @staticmethod\n def _extract_source_types_from_df(df_source: pd.DataFrame) -> set[str]:\n \"\"\"Pull file extensions out of common path/name columns on the source DataFrame.\n\n The direct-upload ingestion path stores extensions in\n ``embedding_metadata.json[source_types]`` so the KB list can render the\n correct file-type icon. When ingestion happens via a connected\n ``input_df`` (e.g. File โ†’ Knowledge) the same field stayed empty and\n the icon defaulted to a blank tile. We look at the well-known columns\n the File / S3 / cloud-storage components produce and collect any\n plausible extension so the icon is consistent across both flows.\n \"\"\"\n candidate_columns = (\"file_path\", \"file_name\", \"filename\", \"source\", \"path\", \"mimetype\")\n extensions: set[str] = set()\n for col in candidate_columns:\n if col not in df_source.columns:\n continue\n for value in df_source[col].dropna():\n ext = KnowledgeComponent._extension_from_value(value)\n # Drop anything that doesn't look like an extension (URL\n # query strings, version segments, etc.) โ€” the icon palette\n # keys off short alphanumeric tokens like \"pdf\"/\"docx\".\n if ext:\n extensions.add(ext)\n return extensions\n\n @staticmethod\n def _extension_from_value(value: Any) -> str | None:\n \"\"\"Return a normalized file-extension token from a path / filename / MIME string.\n\n Accepts values like ``\"report.PDF\"``, ``\"/docs/notes.txt\"``, or\n ``\"application/pdf\"`` and returns ``\"pdf\"`` / ``\"txt\"``. Returns\n ``None`` if no plausible extension can be derived.\n \"\"\"\n extension_length_limit = 10\n if value is None:\n return None\n text = str(value).strip()\n if not text:\n return None\n # MIME types like ``application/pdf`` carry the canonical extension\n # in the subtype slot โ€” preserve them so File / S3 messages keyed\n # only on ``mimetype`` still resolve to an icon.\n if \"/\" in text and \".\" not in text.rsplit(\"/\", 1)[-1]:\n subtype = text.rsplit(\"/\", 1)[-1].strip().lower()\n return subtype if subtype and len(subtype) <= extension_length_limit and subtype.isalnum() else None\n if \".\" not in text:\n return None\n ext = text.rsplit(\".\", 1)[-1].strip().lower()\n if ext and len(ext) <= extension_length_limit and ext.isalnum():\n return ext\n return None\n\n @classmethod\n def _extract_source_types_from_mapping(cls, mapping: Any) -> set[str]:\n \"\"\"Pull extensions from a Message/Data-style ``data`` dict or a plain dict.\"\"\"\n if not isinstance(mapping, dict):\n return set()\n candidate_keys = (\"file_path\", \"file_name\", \"filename\", \"source\", \"path\", \"mimetype\")\n extensions: set[str] = set()\n for key in candidate_keys:\n ext = cls._extension_from_value(mapping.get(key))\n if ext:\n extensions.add(ext)\n return extensions\n\n @classmethod\n def _extract_source_types_from_input(cls, input_value: Any) -> set[str]:\n \"\"\"Pull file extensions out of the raw component input.\n\n ``convert_to_dataframe`` strips Message/Data fields down to ``text``\n when projecting onto a DataFrame, so file metadata attached to a\n File-component \"Raw Content\" output never reaches\n ``_extract_source_types_from_df``. Looking at the raw input first\n keeps the KB icon consistent with direct upload.\n \"\"\"\n if input_value is None:\n return set()\n if isinstance(input_value, list):\n extensions: set[str] = set()\n for item in input_value:\n extensions |= cls._extract_source_types_from_input(item)\n return extensions\n if isinstance(input_value, pd.DataFrame):\n return cls._extract_source_types_from_df(input_value)\n # Message / Data / JSON all expose a ``data`` dict via the lfx schema.\n mapping = getattr(input_value, \"data\", None)\n if mapping is None and isinstance(input_value, dict):\n mapping = input_value\n return cls._extract_source_types_from_mapping(mapping)\n\n def _merge_source_types(self, kb_path: Path, extensions: set[str]) -> None:\n \"\"\"Merge newly observed extensions into the KB's ``source_types`` metadata.\n\n Mirrors the direct-upload path in ``KBIngestionHelper`` so the icon\n rendering on the Knowledge Bases list works regardless of which\n ingestion route was used.\n \"\"\"\n if not extensions:\n return\n metadata_path = kb_path / \"embedding_metadata.json\"\n if not metadata_path.exists():\n return\n try:\n metadata = json.loads(metadata_path.read_text())\n existing = set(metadata.get(\"source_types\") or [])\n metadata[\"source_types\"] = sorted(existing | extensions)\n metadata_path.write_text(json.dumps(metadata, indent=2))\n except (OSError, ValueError, TypeError, json.JSONDecodeError) as e:\n self.log(f\"Warning: Could not update source_types metadata: {e}\")\n\n async def _update_backend_metadata_metrics(self, kb_path: Path, backend: BaseVectorStoreBackend) -> None:\n \"\"\"Update metadata metrics for non-Chroma backends.\"\"\"\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 chunks = await backend.count()\n characters = 0\n words = 0\n async for batch in backend.iter_documents():\n for document in batch:\n characters += len(document.content)\n words += len(document.content.split())\n\n metadata[\"chunks\"] = chunks\n metadata[\"characters\"] = characters\n metadata[\"words\"] = words\n metadata[\"avg_chunk_size\"] = characters / chunks if chunks else 0.0\n metadata[\"size\"] = await backend.storage_size_bytes()\n metadata_path.write_text(json.dumps(metadata, indent=2))\n except (OSError, ValueError, TypeError, json.JSONDecodeError) as e:\n self.log(f\"Warning: Could not update backend metadata metrics: {e}\")\n\n @staticmethod\n def _normalize_backend_selection(value: Any) -> tuple[str, dict[str, Any]]:\n \"\"\"Normalize a DBProviderInput value into backend type/config.\"\"\"\n if not value:\n return BackendType.CHROMA.value, {}\n\n if isinstance(value, str):\n backend_type = value if value == BackendType.OPENSEARCH.value else BackendType.CHROMA.value\n return (\n backend_type,\n _DEFAULT_OPENSEARCH_CONFIG.copy() if backend_type == BackendType.OPENSEARCH.value else {},\n )\n\n if not isinstance(value, dict):\n return BackendType.CHROMA.value, {}\n\n backend_type = str(value.get(\"backend_type\") or value.get(\"id\") or BackendType.CHROMA.value)\n\n if backend_type == BackendType.OPENSEARCH.value:\n backend_config = value.get(\"backend_config\") or value.get(\"config\") or {}\n if not isinstance(backend_config, dict):\n backend_config = {}\n return BackendType.OPENSEARCH.value, {**_DEFAULT_OPENSEARCH_CONFIG, **backend_config}\n\n if backend_type == \"chroma_cloud\":\n backend_config = value.get(\"backend_config\") or value.get(\"config\") or {}\n if not isinstance(backend_config, dict):\n backend_config = {}\n return BackendType.CHROMA.value, {**_DEFAULT_CHROMA_CLOUD_CONFIG, **backend_config}\n\n return BackendType.CHROMA.value, {}\n\n @staticmethod\n def _get_backend_from_metadata(kb_path: Path) -> tuple[str, dict[str, Any]]:\n metadata_path = kb_path / \"embedding_metadata.json\"\n if not metadata_path.exists():\n return BackendType.CHROMA.value, {}\n try:\n metadata = json.loads(metadata_path.read_text())\n except (OSError, json.JSONDecodeError):\n return BackendType.CHROMA.value, {}\n\n backend_type = str(metadata.get(\"backend_type\") or BackendType.CHROMA.value)\n backend_config = metadata.get(\"backend_config\") or {}\n if not isinstance(backend_config, dict):\n backend_config = {}\n return backend_type, backend_config\n\n async def _create_knowledge_base_record(\n self,\n *,\n user_id: Any,\n name: str,\n model_selection: list[dict[str, Any]],\n backend_type: str,\n backend_config: dict[str, Any],\n ) -> None:\n \"\"\"Persist the component-created KB in the DB when Langflow is available.\"\"\"\n try:\n from langflow.api.utils import knowledge_base_service\n except ImportError:\n return\n\n try:\n await knowledge_base_service.create_record(\n user_id=user_id,\n name=name,\n model_selection=model_selection,\n column_config=self.column_config if isinstance(self.column_config, list) else [],\n backend_type=backend_type,\n backend_config=backend_config,\n )\n except Exception as exc: # noqa: BLE001\n self.log(f\"Warning: could not persist knowledge base record: {exc}\")\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 kb_path.mkdir(parents=True, exist_ok=True)\n\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 metadata[\"columns\"].append(\n {\n \"name\": col_name,\n \"vectorize\": vectorize,\n \"identifier\": identifier,\n }\n )\n\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 ) -> BaseVectorStoreBackend:\n \"\"\"Create vector store using the configured DB provider.\"\"\"\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 backend_type, backend_config = self._get_backend_from_metadata(vector_store_dir)\n backend = create_backend(\n backend_type,\n kb_name=self.knowledge_base,\n kb_path=vector_store_dir,\n backend_config=backend_config,\n embedding_function=embedding_function,\n user_id=self.user_id,\n )\n await backend.ensure_ready()\n\n existing_ids = None\n if backend_type != BackendType.CHROMA.value and not self.allow_duplicates:\n existing_ids = set()\n async for batch in backend.iter_documents():\n for document in batch:\n doc_id = document.metadata.get(\"_id\")\n if doc_id:\n existing_ids.add(doc_id)\n\n data_objects = await self._convert_df_to_data_objects(df_source, config_list, existing_ids=existing_ids)\n\n user_metadata_tag = self._resolve_user_metadata_tag()\n\n documents = []\n for data_obj in data_objects:\n doc = data_obj.to_lc_document()\n if user_metadata_tag:\n doc.metadata[\"source_metadata\"] = user_metadata_tag\n documents.append(doc)\n\n if documents:\n await backend.add_documents(documents)\n self.log(f\"Added {len(documents)} documents to vector store '{self.knowledge_base}'\")\n\n return backend\n\n async def _convert_df_to_data_objects(\n self,\n df_source: pd.DataFrame,\n config_list: list[dict[str, Any]],\n existing_ids: set[str] | None = None,\n ) -> list[Data]:\n \"\"\"Convert DataFrame to Data objects for vector store.\"\"\"\n data_objects: list[Data] = []\n\n if existing_ids is None:\n kb_path = await self._kb_path()\n\n chroma = Chroma(\n persist_directory=str(kb_path),\n collection_name=self.knowledge_base,\n **chroma_langchain_collection_kwargs(),\n )\n\n all_docs = chroma.get()\n\n existing_ids = {metadata.get(\"_id\") for metadata in all_docs[\"metadatas\"] if metadata.get(\"_id\")}\n\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 if identifier:\n identifier_cols.append(col_name)\n\n for _, row in df_source.iterrows():\n identifier_parts = [str(row[col]) for col in content_cols if col in row and self._scalar_notna(row[col])]\n\n page_content = \" \".join(identifier_parts)\n\n data_dict = {\n \"text\": page_content,\n }\n\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 for col in df_source.columns:\n if col not in content_cols and col in row and self._scalar_notna(row[col]):\n value = row[col]\n data_dict[col] = str(value)\n\n page_content_hash = hashlib.sha256(page_content.encode()).hexdigest()\n data_dict[\"_id\"] = page_content_hash\n\n if not self.allow_duplicates and page_content_hash in existing_ids:\n self.log(f\"Skipping duplicate row with hash {page_content_hash}\")\n continue\n\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 \"\"\"Validate collection name.\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 if not (min_length <= len(name) <= max_length):\n return False\n\n if not (name[0].isalnum() and name[-1].isalnum()):\n return False\n\n return re.match(r\"^[a-zA-Z0-9_-]+$\", name) is not None\n\n async def _kb_path(self) -> Path | None:\n cached_path = getattr(self, \"_cached_kb_path\", None)\n if cached_path is not None:\n return cached_path\n\n # Lazy import to keep ``lfx`` importable standalone โ€” langflow's\n # user/DB models are not always available at module load time.\n from langflow.services.database.models.user.crud import get_user_by_id\n\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 self._cached_kb_path = kb_root / kb_user / self.knowledge_base\n\n return self._cached_kb_path\n\n def _resolve_user_metadata_tag(self) -> str:\n \"\"\"Return the JSON-encoded user metadata tag for chunk writes.\"\"\"\n raw = getattr(self, \"metadata_json\", None)\n if not raw:\n return \"\"\n text = raw.strip() if isinstance(raw, str) else raw\n if not text:\n return \"\"\n try:\n decoded = json.loads(text)\n except (TypeError, json.JSONDecodeError) as exc:\n self.log(f\"KnowledgeComponent: metadata_json is not valid JSON ({exc}); skipping metadata stamp.\")\n return \"\"\n if not isinstance(decoded, dict):\n self.log(\"KnowledgeComponent: metadata_json must decode to a JSON object; skipping metadata stamp.\")\n return \"\"\n return json.dumps(decoded, sort_keys=True)\n\n async def build_kb_info(self) -> Data:\n \"\"\"Main ingestion routine โ†’ returns a dict with KB metadata.\n\n The annotation is intentionally narrowed to ``Data`` even though the\n cross-mode fallback below may return a ``DataFrame`` from\n ``retrieve_data``. The frontend builds React-Flow handle IDs from\n this output's type list; widening it to ``Data | DataFrame`` makes\n the API advertise ``[\"JSON\", \"Table\"]`` for the ingest output and\n breaks every saved-edge sourceHandle that was generated against\n a single-type handle (BUG-02). Python doesn't enforce return\n annotations at runtime, so the rare fallback path keeps working.\n \"\"\"\n if _is_retrieve_mode(getattr(self, \"mode\", MODE_INGEST)):\n return await self.retrieve_data()\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n\n run_id: uuid.UUID | None = None\n run_job_id: uuid.UUID | None = None\n run_summary: IngestionSummary | None = None\n run_status: IngestionRunStatus = IngestionRunStatus.SUCCEEDED\n run_error: str | None = None\n kb_record_id: uuid.UUID | None = None\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 config_list = self._validate_column_config(df_source)\n column_metadata = self._build_column_metadata(config_list, df_source)\n\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 if metadata_path.exists():\n settings_service = get_settings_service()\n stored_metadata = json.loads(metadata_path.read_text())\n\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 embedding_model_name = stored_metadata.get(\"embedding_model\")\n embedding_provider = stored_metadata.get(\"embedding_provider\", \"Unknown\")\n if embedding_model_name:\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 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 if not self.api_key:\n log_label = f\"knowledge base '{self.knowledge_base}'\"\n msg = (\n f\"Cannot decrypt the stored embedding API key for {log_label}. \"\n \"This usually means the server's SECRET_KEY changed after the \"\n \"key was saved. To recover, supply the embedding provider API \"\n \"key on the component's 'Embedding Provider API Key' input and \"\n \"re-run ingestion โ€” the key will be re-encrypted with the \"\n \"current SECRET_KEY.\"\n )\n raise KBKeyDecryptError(msg) from e\n logger.warning(\"Stored API key undecryptable; using component-supplied key. Error: %s\", e)\n\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 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 run_id, run_job_id, run_summary, kb_record_id = await self._begin_ingestion_run(kb_path)\n if kb_record_id is not None:\n await self._record_kb_status(kb_record_id, \"ingesting\")\n\n backend = await self._create_vector_store(df_source, config_list, embedding_function=embedding_function)\n\n self._save_kb_files(kb_path, config_list)\n\n try:\n if not isinstance(backend, BaseVectorStoreBackend):\n pass\n elif backend.backend_type == BackendType.CHROMA and hasattr(backend, \"raw_langchain_store\"):\n self._update_metadata_metrics(kb_path, backend.raw_langchain_store())\n else:\n await self._update_backend_metadata_metrics(kb_path, backend)\n # Stamp the KB with the file extensions we just ingested so\n # the Knowledge Bases list renders the correct icon for\n # flow-driven ingestion (input_df), matching direct upload.\n # We look at the raw input first because ``convert_to_dataframe``\n # drops Message/Data metadata fields (file_path, mimetype, โ€ฆ)\n # when projecting onto the DataFrame.\n source_types = self._extract_source_types_from_input(input_value)\n source_types |= self._extract_source_types_from_df(df_source)\n self._merge_source_types(kb_path, source_types)\n finally:\n if isinstance(backend, BaseVectorStoreBackend):\n await backend.teardown()\n\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 if run_summary is not None:\n run_summary.record_item(\n IngestionItemResult(\n item_id=self.knowledge_base,\n display_name=f\"{self.knowledge_base} ({len(df_source)} rows)\",\n status=IngestionItemStatus.SUCCEEDED,\n chunks_created=len(df_source),\n ),\n )\n\n if kb_record_id is not None:\n await self._record_kb_stats(kb_record_id, kb_path)\n await self._record_kb_status(kb_record_id, \"ready\")\n\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 run_status = IngestionRunStatus.FAILED\n run_error = str(e) or e.__class__.__name__\n msg = f\"Error during KB ingestion: {e}\"\n raise RuntimeError(msg) from e\n except Exception as e:\n run_status = IngestionRunStatus.FAILED\n run_error = str(e) or e.__class__.__name__\n raise\n finally:\n if run_id is not None and run_summary is not None:\n await self._finalize_ingestion_run(\n run_id=run_id,\n job_id=run_job_id,\n summary=run_summary,\n status=run_status,\n error_message=run_error,\n )\n if kb_record_id is not None and run_status is IngestionRunStatus.FAILED:\n await self._record_kb_status(kb_record_id, \"failed\", failure_reason=run_error)\n\n async def _begin_ingestion_run(\n self,\n kb_path: Path,\n ) -> tuple[uuid.UUID | None, uuid.UUID | None, IngestionSummary | None, uuid.UUID | None]:\n \"\"\"Create a parent ``Job`` and seed an ingestion-run row.\"\"\"\n if not self.user_id:\n self.log(\"No user_id on component; skipping ingestion-run tracking.\")\n return None, None, None, None\n\n try:\n user_uuid = uuid.UUID(str(self.user_id))\n except (ValueError, TypeError) as exc:\n self.log(f\"Could not coerce user_id={self.user_id!r} to UUID; skipping run tracking ({exc}).\")\n return None, None, None, None\n\n try:\n from langflow.api.utils import ingestion_run_service, knowledge_base_service\n from langflow.services.database.models.jobs.model import JobStatus, JobType\n from langflow.services.deps import get_job_service\n except ImportError as exc:\n self.log(f\"Run-history wiring unavailable; ingestion will not be recorded ({exc}).\")\n return None, None, None, None\n\n try:\n kb_record = await knowledge_base_service.get_by_user_and_name(user_uuid, self.knowledge_base)\n kb_record_id = kb_record.id if kb_record is not None else None\n\n user_metadata = self._parse_user_metadata_dict()\n\n job_id = uuid.uuid4()\n job_service = get_job_service()\n raw_flow_id = getattr(self, \"flow_id\", None)\n flow_id_uuid = uuid.UUID(str(raw_flow_id)) if raw_flow_id else job_id\n await job_service.create_job(\n job_id=job_id,\n flow_id=flow_id_uuid,\n job_type=JobType.INGESTION,\n asset_id=kb_record_id,\n asset_type=\"knowledge_base\",\n user_id=user_uuid,\n )\n await job_service.update_job_status(job_id, JobStatus.IN_PROGRESS)\n\n source = FlowComponentSource(\n user_id=user_uuid,\n source_config={\n \"knowledge_base\": self.knowledge_base,\n \"kb_path\": str(kb_path),\n \"flow_id\": str(flow_id_uuid),\n },\n )\n\n run_id = await ingestion_run_service.create_run(\n kb_name=self.knowledge_base,\n source=source,\n job_id=job_id,\n user_id=user_uuid,\n kb_id=kb_record_id,\n user_metadata=user_metadata,\n )\n if run_id is None:\n return None, job_id, None, kb_record_id\n await ingestion_run_service.mark_running(run_id)\n\n summary = IngestionSummary(\n kb_name=self.knowledge_base,\n source_type=source.source_type.value,\n user_id=user_uuid,\n job_id=job_id,\n source_config=source.describe().get(\"config\") or {},\n user_metadata=user_metadata,\n )\n self.log(f\"Started ingestion run job_id={job_id} kb_name={self.knowledge_base} kb_id={kb_record_id}\")\n except Exception as exc: # noqa: BLE001 โ€” telemetry must never abort ingestion\n self.log(f\"Could not begin ingestion-run tracking: {exc}\")\n return None, None, None, None\n else:\n return run_id, job_id, summary, kb_record_id\n\n async def _finalize_ingestion_run(\n self,\n *,\n run_id: uuid.UUID,\n job_id: uuid.UUID | None,\n summary: IngestionSummary,\n status: IngestionRunStatus,\n error_message: str | None,\n ) -> None:\n \"\"\"Persist the final summary and transition the parent Job.\"\"\"\n try:\n from langflow.api.utils import ingestion_run_service\n from langflow.services.database.models.jobs.model import JobStatus\n from langflow.services.deps import get_job_service\n except ImportError as exc:\n self.log(f\"Run-history wiring unavailable; ingestion-run finalize skipped ({exc}).\")\n return\n\n try:\n await ingestion_run_service.finalize_run(\n run_id,\n summary=summary,\n status=status,\n error_message=error_message,\n )\n if job_id is not None:\n terminal_status = JobStatus.COMPLETED if status is not IngestionRunStatus.FAILED else JobStatus.FAILED\n await get_job_service().update_job_status(job_id, terminal_status, finished_timestamp=True)\n except Exception as exc: # noqa: BLE001 โ€” telemetry must never re-raise\n self.log(f\"Could not finalize ingestion-run tracking: {exc}\")\n\n async def _record_kb_status(\n self,\n kb_record_id: uuid.UUID,\n status_value: str,\n *,\n failure_reason: str | None = None,\n ) -> None:\n \"\"\"Mirror Path A's KB-row status transitions.\"\"\"\n try:\n from langflow.api.utils import knowledge_base_service\n from langflow.services.database.models.knowledge_base.model import KnowledgeBaseStatus\n except ImportError:\n return\n try:\n await knowledge_base_service.update_status(\n kb_record_id,\n status=KnowledgeBaseStatus(status_value),\n failure_reason=failure_reason,\n )\n except Exception as exc: # noqa: BLE001\n self.log(f\"Could not update KB status to {status_value}: {exc}\")\n\n async def _record_kb_stats(self, kb_record_id: uuid.UUID, kb_path: Path) -> None:\n \"\"\"Push freshly-refreshed metrics from embedding_metadata.json onto the DB row.\"\"\"\n try:\n from langflow.api.utils import knowledge_base_service\n except ImportError:\n self.log(\"knowledge_base_service unavailable; KB stats will not sync to DB row.\")\n return\n metadata_path = kb_path / \"embedding_metadata.json\"\n if not metadata_path.exists():\n self.log(f\"No embedding_metadata.json at {metadata_path}; skipping KB stats sync.\")\n return\n try:\n metadata = json.loads(metadata_path.read_text())\n except (OSError, json.JSONDecodeError) as exc:\n self.log(f\"Could not read KB metadata for stats sync: {exc}\")\n return\n chunks = int(metadata.get(\"chunks\", 0) or 0)\n words = int(metadata.get(\"words\", 0) or 0)\n characters = int(metadata.get(\"characters\", 0) or 0)\n try:\n await knowledge_base_service.update_stats(\n kb_record_id,\n chunks=chunks,\n words=words,\n characters=characters,\n size_bytes=int(metadata.get(\"size\", 0) or 0),\n source_types=list(metadata.get(\"source_types\") or []),\n chunk_size=metadata.get(\"chunk_size\"),\n chunk_overlap=metadata.get(\"chunk_overlap\"),\n separator=metadata.get(\"separator\"),\n )\n self.log(f\"Synced KB stats to DB row {kb_record_id}: chunks={chunks} words={words} characters={characters}\")\n except Exception as exc: # noqa: BLE001\n self.log(f\"Could not sync KB stats to DB row {kb_record_id}: {exc}\")\n\n def _parse_user_metadata_dict(self) -> dict[str, Any]:\n \"\"\"Decode ``metadata_json`` to a dict, or ``{}`` on any error.\"\"\"\n raw = getattr(self, \"metadata_json\", None)\n if not raw:\n return {}\n text = raw.strip() if isinstance(raw, str) else raw\n if not text:\n return {}\n try:\n decoded = json.loads(text)\n except (TypeError, json.JSONDecodeError):\n return {}\n return decoded if isinstance(decoded, dict) else {}\n\n # =====================================================================\n # RETRIEVAL CODE PATH\n # =====================================================================\n @property\n def _user_uuid(self) -> uuid.UUID | None:\n \"\"\"Return self.user_id as a UUID, converting from str if necessary.\"\"\"\n if not self.user_id:\n return None\n return self.user_id if isinstance(self.user_id, uuid.UUID) else uuid.UUID(self.user_id)\n\n def _get_kb_metadata(self, kb_path: Path, *, require_api_key: bool = False) -> dict:\n \"\"\"Load the knowledge base's embedding metadata file.\"\"\"\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n return load_kb_metadata(\n kb_path,\n log_label=f\"knowledge base '{self.knowledge_base}'\",\n require_api_key=require_api_key,\n )\n\n async def _resolve_backend(self, *, kb_user: str) -> tuple[str, dict[str, Any]]: # noqa: ARG002\n \"\"\"Return ``(backend_type, backend_config)`` for this KB.\"\"\"\n try:\n from langflow.api.utils import knowledge_base_service\n\n user_uuid = self._user_uuid\n if user_uuid is None:\n return BackendType.CHROMA.value, {}\n record = await knowledge_base_service.get_by_user_and_name(user_uuid, self.knowledge_base)\n except Exception as exc: # noqa: BLE001\n logger.debug(\"KB record lookup failed: %s\", exc)\n return BackendType.CHROMA.value, {}\n\n if record is None:\n return BackendType.CHROMA.value, {}\n return (\n record.backend_type or BackendType.CHROMA.value,\n record.backend_config or {},\n )\n\n def _resolve_model_selection(self, metadata: dict[str, Any]) -> list[dict[str, Any]]:\n \"\"\"Resolve the ``get_embeddings``-compatible model selection from metadata.\"\"\"\n model_selection = metadata.get(\"model_selection\")\n if model_selection:\n selection_list = [model_selection] if isinstance(model_selection, dict) else list(model_selection)\n return [self._hydrate_model_metadata(entry) for entry in selection_list]\n\n embedding_model_name = metadata.get(\"embedding_model\")\n embedding_provider = metadata.get(\"embedding_provider\", \"Unknown\")\n if not embedding_model_name:\n msg = (\n f\"Knowledge base '{self.knowledge_base}' has no embedding model recorded; \"\n \"re-create it with a supported embedding model.\"\n )\n raise ValueError(msg)\n\n match = self._find_catalog_entry(embedding_model_name)\n if match is None:\n msg = (\n f\"Embedding model '{embedding_model_name}' (provider '{embedding_provider}') \"\n \"recorded for this knowledge base is no longer available in the model registry. \"\n \"Please re-create the knowledge base with a supported embedding model.\"\n )\n raise ValueError(msg)\n return [match]\n\n def _hydrate_model_metadata(self, entry: dict[str, Any]) -> dict[str, Any]:\n \"\"\"Fill in ``metadata.embedding_class`` / ``param_mapping`` if missing.\"\"\"\n entry_metadata = entry.get(\"metadata\") or {}\n has_class = bool(entry_metadata.get(\"embedding_class\"))\n has_mapping = bool(entry_metadata.get(\"param_mapping\"))\n if has_class and has_mapping:\n return entry\n\n model_name = entry.get(\"name\")\n if not model_name:\n return entry\n\n catalog_entry = self._find_catalog_entry(model_name)\n if catalog_entry is None:\n return entry\n\n catalog_metadata = catalog_entry.get(\"metadata\") or {}\n merged_metadata = {**catalog_metadata, **entry_metadata}\n return {**entry, \"metadata\": merged_metadata}\n\n def _find_catalog_entry(self, model_name: str) -> dict[str, Any] | None:\n \"\"\"Look up an embedding model by name in the unified-models catalog.\"\"\"\n options = get_embedding_model_options(user_id=self.user_id)\n return next((o for o in options if o.get(\"name\") == model_name), None)\n\n async def retrieve_data(self) -> DataFrame:\n \"\"\"Retrieve data from the selected knowledge base.\n\n Annotation narrowed to ``DataFrame`` to keep this output's handle\n type list at ``[\"Table\"]`` only; widening to a union surfaces\n ``[\"Table\", \"JSON\"]`` on the API and breaks every starter-project\n edge whose sourceHandle was stored against a single-type handle\n (BUG-02). The cross-mode defensive fallback may still return a\n ``Data`` at runtime โ€” Python ignores return annotations, so\n nothing breaks.\n \"\"\"\n if not _is_retrieve_mode(getattr(self, \"mode\", MODE_INGEST)):\n return await self.build_kb_info()\n raise_error_if_astra_cloud_disable_component(astra_error_msg)\n\n # Lazy import: langflow's user/DB models aren't part of lfx's\n # standalone install, so ``lfx run .json`` can't resolve\n # this symbol at module import time. Deferring to use keeps the\n # component importable in both environments.\n from langflow.services.database.models.user.crud import get_user_by_id\n\n async with session_scope() as db:\n if not self.user_id:\n msg = \"User ID is required for fetching Knowledge Base data.\"\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 kb_path = _get_knowledge_bases_root_path() / kb_user / self.knowledge_base\n\n component_api_key = self.api_key if getattr(self, \"api_key\", None) else None\n needs_stored_key = not component_api_key\n metadata = self._get_kb_metadata(kb_path, require_api_key=needs_stored_key)\n if not metadata:\n msg = f\"Metadata not found for knowledge base: {self.knowledge_base}. Ensure it has been indexed.\"\n raise ValueError(msg)\n\n api_key = component_api_key or metadata.get(\"api_key\")\n model_selection = self._resolve_model_selection(metadata)\n chunk_size = metadata.get(\"chunk_size\")\n embedding_function = get_embeddings(\n model=model_selection,\n user_id=self.user_id,\n api_key=api_key,\n chunk_size=chunk_size,\n )\n\n backend_type, backend_config = await self._resolve_backend(kb_user=kb_user)\n backend = create_backend(\n backend_type,\n kb_name=self.knowledge_base,\n kb_path=kb_path,\n backend_config=backend_config,\n embedding_function=embedding_function,\n user_id=self.user_id,\n )\n try:\n user_metadata_filter = _parse_metadata_filter(getattr(self, \"metadata_filter\", None))\n use_scores = bool(self.search_query)\n search_k = self.top_k * 4 if user_metadata_filter else self.top_k\n results = await backend.similarity_search(\n query=self.search_query or \"\",\n k=search_k,\n with_scores=use_scores,\n )\n if user_metadata_filter:\n results = [\n (doc, score) for doc, score in results if _chunk_matches_filter(doc.metadata, user_metadata_filter)\n ]\n results = results[: self.top_k]\n\n embeddings_by_key: dict[tuple[str, str], list[float]] = {}\n if self.include_embeddings and results:\n # Join each retrieved chunk to its stored embedding. We key on\n # ``_id`` when present and fall back to page content otherwise,\n # so KBs populated by direct file upload โ€” whose chunks carry no\n # ``_id`` โ€” still resolve. See ``_embedding_match_key``.\n wanted_keys = {_embedding_match_key(doc.page_content, doc.metadata) for doc, _score in results}\n async for batch in backend.iter_documents(include_embeddings=True):\n for entry in batch:\n if entry.embedding is None:\n continue\n key = _embedding_match_key(entry.content, entry.metadata)\n if key in wanted_keys:\n embeddings_by_key[key] = entry.embedding\n if len(embeddings_by_key) == len(wanted_keys):\n break\n\n data_list: list[Data] = []\n for doc, score in results:\n kwargs: dict[str, Any] = {\"content\": doc.page_content}\n if use_scores:\n kwargs[\"_score\"] = -1 * score\n if self.include_metadata:\n kwargs.update(doc.metadata)\n if self.include_embeddings:\n kwargs[\"_embeddings\"] = embeddings_by_key.get(_embedding_match_key(doc.page_content, doc.metadata))\n data_list.append(Data(**kwargs))\n\n return DataFrame(data=data_list)\n finally:\n await backend.teardown()\n\n\ndef _embedding_match_key(content: str, metadata: dict[str, Any] | None) -> tuple[str, str]:\n \"\"\"Build a stable key for aligning a retrieved chunk with its stored embedding.\n\n Embeddings are gathered separately (via ``iter_documents``) and then joined\n back onto the search results. Component-driven ingestion stamps a\n content-hash ``_id`` on every chunk, but direct file-upload ingestion\n (``KBIngestionHelper.perform_ingestion``) does not. Keying the join purely on\n ``_id`` therefore left every upload-populated KB with ``_embeddings: None``.\n\n We prefer ``_id`` when present (so legitimately distinct chunks that happen to\n share text aren't collapsed) and fall back to the chunk's page content\n otherwise. The content fallback is exact for the embedding use case: two\n chunks with identical text necessarily share the same embedding vector. The\n leading ``\"id\"`` / ``\"content\"`` tag namespaces the two key spaces so a mixed\n KB (some chunks with ``_id``, some without) never cross-matches.\n \"\"\"\n doc_id = metadata.get(\"_id\") if metadata else None\n if doc_id:\n return (\"id\", str(doc_id))\n return (\"content\", content or \"\")\n\n\ndef _parse_metadata_filter(raw: str | None) -> dict[str, list[str]]:\n \"\"\"Decode the ``metadata_filter`` input into a {key: [values]} map.\n\n Empty or malformed input maps to an empty filter so retrieval falls back\n to the unfiltered path. JSON errors are swallowed here rather than raised:\n surfacing component-config errors at the canvas node would break a flow\n run for what is meant to be an optional refinement.\n \"\"\"\n if not raw:\n return {}\n text = raw.strip() if isinstance(raw, str) else raw\n if not text:\n return {}\n try:\n decoded = json.loads(text)\n except (TypeError, json.JSONDecodeError):\n logger.warning(\"KnowledgeComponent: metadata_filter is not valid JSON; ignoring filter.\")\n return {}\n if not isinstance(decoded, dict):\n logger.warning(\"KnowledgeComponent: metadata_filter must be a JSON object; ignoring filter.\")\n return {}\n result: dict[str, list[str]] = {}\n for key, value in decoded.items():\n if not isinstance(key, str):\n continue\n if isinstance(value, list):\n result[key] = [str(entry) for entry in value]\n else:\n result[key] = [str(value)]\n return result\n\n\ndef _chunk_matches_filter(metadata: dict[str, Any] | None, filt: dict[str, list[str]]) -> bool:\n \"\"\"AND across keys, OR within key values, mirroring the chunks endpoint.\"\"\"\n if not filt:\n return True\n if not metadata:\n return False\n raw = metadata.get(\"source_metadata\")\n if not raw:\n return False\n try:\n stored = json.loads(raw) if isinstance(raw, str) else raw\n except json.JSONDecodeError:\n return False\n if not isinstance(stored, dict):\n return False\n for key, expected_values in filt.items():\n actual = stored.get(key)\n if actual is None:\n return False\n actual_set = {str(entry) for entry in actual} if isinstance(actual, list) else {str(actual)}\n if not actual_set & set(expected_values):\n return False\n return True\n" }, "column_config": { "_input_type": "TableInput", @@ -30402,6 +30402,6 @@ "num_components": 129, "num_modules": 17 }, - "sha256": "82ee44f7bed230621fc56ceb61d935bbc2f462bbcb56e1d7d408f5ca130d60db", + "sha256": "420e8f624e9e8b9f8641471c89f1b20f5512e0e3cd3f683b1c19db2f2fcc0cac", "version": "1.11.0" }