feat: Add support for Ingestion and Retrieval of Knowledge Bases (#9088)
* refactor: Standardize import statements and improve code readability across components - Updated import statements to use consistent single quotes. - Refactored various components to enhance readability and maintainability. - Adjusted folder and file handling logic in the sidebar and file manager components. - Introduced a new tabbed interface for the files page to separate files and knowledge bases, improving user experience. * [autofix.ci] apply automated fixes * feat: Introduce new Files and Knowledge Bases page with tabbed interface - Added a new FilesPage component to manage file uploads and organization. - Implemented a tabbed interface to separate Files and Knowledge Bases for improved user experience. - Created FilesTab and KnowledgeBasesTab components for handling respective functionalities. - Refactored routing to accommodate the new structure and updated import statements for consistency. - Removed the old filesPage component to streamline the codebase. * Create knowledgebase_utils.py * Push initial ingest component * [autofix.ci] apply automated fixes * Create initial KB Ingestion component * [autofix.ci] apply automated fixes * Fix ruff check on utility functions * [autofix.ci] apply automated fixes * Some quick fixes * Update kb_ingest.py * [autofix.ci] apply automated fixes * First version of retrieval component * [autofix.ci] apply automated fixes * Update icon * Update kb_retrieval.py * [autofix.ci] apply automated fixes * Add knowledge bases feature with API integration and UI components * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * Refactor imports and update routing paths for assets and main page components. Adjust tab handling in the assets page to reflect URL changes and improve user navigation experience. * [autofix.ci] apply automated fixes * Add CreateKnowledgeBaseButton, KnowledgeBaseEmptyState, and KnowledgeBaseSelectionOverlay components. Refactor KnowledgeBasesTab to utilize new components and improve UI for knowledge base management. Introduce utility functions for formatting numbers and average chunk sizes. * [autofix.ci] apply automated fixes * PoV: Add Parquet data retrieval to KBRetrievalComponent (#9097) * Add Parquet data retrieval to KBRetrievalComponent Introduces a new output to KBRetrievalComponent for returning knowledge base data by reading Parquet files. Updates dependencies to include fastparquet for Parquet support. * [autofix.ci] apply automated fixes --------- Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com> * Fix some ruff issues * [autofix.ci] apply automated fixes * feat: refactor file management and knowledge base components - Replaced the existing assetsPage with a new filesPage to better organize file management functionalities. - Introduced KnowledgePage to handle knowledge base operations, integrating KnowledgeBasesTab for displaying and managing knowledge bases. - Added various components for file and knowledge base management, including CreateKnowledgeBaseButton, KnowledgeBaseEmptyState, and drag-and-drop functionality. - Updated routing and imports to reflect the new structure and ensure consistency across the application. - Enhanced user experience with improved UI elements and state management for file selection and operations. * feat: implement delete confirmation modal for knowledge base deletion - Added a DeleteConfirmationModal component to confirm deletion actions. - Integrated the modal into the KnowledgeBasesTab for handling knowledge base deletions. - Updated column definitions to include a delete button for each knowledge base. - Enhanced user experience by ensuring deletion actions require confirmation. - Adjusted styles for the knowledge base table to improve checkbox visibility. * feat: enhance knowledge base metadata with embedding model detection - Added `embedding_model` field to `KnowledgeBaseInfo` for improved metadata tracking. - Implemented `detect_embedding_model` function to extract embedding model information from configuration files. - Updated `get_kb_metadata` to prioritize metadata extraction from `embedding_metadata.json`, falling back to detection if necessary. - Modified `KBIngestionComponent` to save embedding model metadata during ingestion. - Adjusted frontend components to display embedding model information in knowledge base queries and tables. * refactor: clean up tooltip and value getter comments in knowledge base columns - Removed redundant comments in the `knowledgeBaseColumns.tsx` file to enhance code clarity. - Simplified the tooltip and value getter functions for embedding model display. * [autofix.ci] apply automated fixes * refactor: simplify KnowledgeBaseSelectionOverlay component - Removed the unused onExport prop and its associated functionality. - Cleaned up code formatting for consistency and readability. - Updated success message strings to use single quotes for uniformity. * feat: implement bulk and single deletion for knowledge bases - Added `BulkDeleteRequest` model to handle bulk deletion requests. - Implemented `delete_knowledge_base` endpoint for single knowledge base deletion. - Created `delete_knowledge_bases_bulk` endpoint for deleting multiple knowledge bases at once. - Introduced `useDeleteKnowledgeBase` and `useDeleteKnowledgeBases` hooks for frontend integration. - Updated `KnowledgeBaseSelectionOverlay` and `KnowledgeBasesTab` components to utilize new deletion functionality with user feedback on success and error handling. * Initial support for vector search * feat: add KnowledgeBaseDrawer component for enhanced knowledge base details - Introduced `KnowledgeBaseDrawer` component to display detailed information about selected knowledge bases. - Integrated mock data for source files and linked flows, with a layout for displaying descriptions and embedding models. - Updated `KnowledgeBasesTab` to handle row clicks and open the drawer with relevant knowledge base data. - Enhanced `KnowledgePage` to manage drawer state and selected knowledge base, improving user interaction and experience. * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * Fix ruff checks * Update knowledge_bases.py * feat: update mock data and enhance drawer functionality in KnowledgeBase components - Replaced mock data in `KnowledgeBaseDrawer` with more descriptive placeholders. - Added a reference to the drawer in `KnowledgePage` for improved click handling. - Implemented logic to close the drawer when clicking outside, except for table row clicks. - Enhanced row click handling to toggle drawer state based on current visibility. * [autofix.ci] apply automated fixes * Append scores column to rows * refactor: improve knowledge base deletion and UI components - Updated `useDeleteKnowledgeBase` and `useDeleteKnowledgeBases` to enhance parameter naming for clarity. - Removed the `CreateKnowledgeBaseButton` component and its references to streamline the UI. - Simplified the `KnowledgeBaseDrawer` and `KnowledgeBasesTab` components by removing mock data and improving state management. - Enhanced the `KnowledgeBaseSelectionOverlay` to better handle bulk deletions and selection states. - Refactored various components for consistent styling and improved readability. * refactor: standardize import statements and improve code readability in SideBarFoldersButtonsComponent - Updated import statements to use consistent single quotes. - Refactored various function calls and state management for improved clarity. - Enhanced folder handling logic and UI interactions for better user experience. * feat: Add encryption for API keys in KB ingest and retrieval (#9129) Add encryption for API keys in KB ingest and retrieval Introduces secure storage of embedding model API keys by encrypting them during knowledge base ingestion and decrypting them during retrieval. Refactors metadata handling to include encrypted API keys, updates retrieval to support decryption and dynamic embedder construction, and improves logging for key operations. Removes legacy embedding client code in retrieval in favor of a provider-based approach. * [autofix.ci] apply automated fixes * Fix import of auth utils * Allow appending to existing knowledge base * [autofix.ci] apply automated fixes * Update kb_ingest.py * Update kb_ingest.py * feat: enhance table component with editable Vectorize column functionality - Implemented logic to determine editability of the Vectorize column based on other row values. - Added checks to refresh grid cells upon changes to the Vectorize column. - Updated TableAutoCellRender to conditionally disable editing based on Vectorize column state. * New ingestion creation dialog * [autofix.ci] apply automated fixes * Clean up the creation process for KB * [autofix.ci] apply automated fixes * Clean up names and descriptions * Update kb_retrieval.py * chroma retrieval * [autofix.ci] apply automated fixes * Further KB cleanup * refactor: update KB ingestion component and enhance NodeDialog functionality - Restored SecretStrInput for API key in KB ingestion component. - Modified NodeDialog to handle new value format and added support for additional properties. - Introduced custom hooks for managing global variable states in InputGlobalComponent. - Improved dropdown component styling and interaction. - Cleaned up input component code for better readability and maintainability. * Hash the text as id * [autofix.ci] apply automated fixes * Update kb_retrieval.py * [autofix.ci] apply automated fixes * Make sure to write out the source parquet * Remove unneeded old code * Add ability to block duplicate ingestion chunks * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * Rename retrieval component * Better refresh mechanism for the retrieve * Clean up some unused functionality * Update kb_ingest.py * Fix dropdown component logic to include checks for refresh button and dialog inputs * Test the API key before saving knowledge * [autofix.ci] apply automated fixes * Allow storing updated api keys if provided at ingest time * Add Knowledge Bases component and enhance Knowledge Base Empty State - Introduced a new JSON configuration for Knowledge Bases, defining nodes and edges for data processing. - Enhanced the KnowledgeBaseEmptyState component to include a button for creating a knowledge base template. - Updated KnowledgeBasesTab to handle template creation, integrating flow management and navigation features. * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * Update Knowledge Bases.json * Update Knowledge Bases configuration and enhance UI components - Updated the code hash in the Knowledge Bases JSON configuration. - Modified the KnowledgeBaseEmptyState component to change the button icon and text from "Try Knowledge Base Template" to "Create Knowledge". - Cleared the options for the knowledge base selection dropdowns to ensure they reflect the current state of available knowledge bases. * [autofix.ci] apply automated fixes * Implement feature flag for Knowledge Bases functionality - Added FEATURE_FLAGS.knowledge_bases to control the visibility of knowledge base components in the API and UI. - Updated the router to conditionally include the knowledge bases router based on the feature flag. - Modified KBIngestionComponent and KBRetrievalComponent to hide if the knowledge bases feature is disabled. - Enhanced the initial setup to skip loading knowledge base starter projects when the feature is disabled. - Updated frontend routes and sidebar components to conditionally render knowledge base options based on the feature flag. - Adjusted API queries to return an empty array if the knowledge bases feature is disabled. * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * Refactor Knowledge Bases feature flag implementation - Removed the FEATURE_FLAGS.knowledge_bases flag from backend components and frontend routes. - Updated the API and UI to always include knowledge base components, simplifying the codebase. - Adjusted the frontend feature flags to set ENABLE_KNOWLEDGE_BASES to false, ensuring knowledge base features are not displayed. - Cleaned up related components and routes to reflect the removal of the feature flag, enhancing maintainability. * revert * [autofix.ci] apply automated fixes * Remove Knowledge Bases JSON configuration and clean up KnowledgeBasesTab component by eliminating unused imports and template creation functionality. * [autofix.ci] apply automated fixes * Enhance routing structure by adding admin and login routes with protected access. Refactor flow routes for improved organization and clarity. * added template back * Use chroma for stats computation * Fix ruff issue * [autofix.ci] apply automated fixes * Update Knowledge Bases.json * Update Knowledge Bases.json * Rename to just knowledge * feat: enhance Jest configuration and add new tests for Knowledge Base components - Updated jest.config.js to include a new setup file and refined test matching patterns. - Introduced jest.setup.js for mocking globals and Vite-specific syntax. - Added tests for KnowledgeBaseDrawer, KnowledgeBaseEmptyState, KnowledgeBaseSelectionOverlay, KnowledgeBasesTab, and KnowledgePage components. - Created utility functions for testing and mock data for knowledge bases. - Implemented tests for utility functions related to knowledge base formatting. * [autofix.ci] apply automated fixes * refactor: reorganize imports and clean up console log in Dropdown component - Moved and re-imported necessary dependencies for better structure. - Removed unnecessary console log statement to clean up the code. * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * feat: add success callback for knowledge base creation in NodeDialog component - Introduced a new success callback to handle knowledge base creation notifications. - Enhanced dialog closing logic with a delay for Astra database tracking. - Reorganized imports for better structure. * refactor: update table component to handle single-toggle columns - Renamed functions and variables to improve clarity regarding single-toggle columns (Vectorize and Identifier). - Updated logic to ensure proper editability checks for single-toggle columns. - Adjusted related components to reflect changes in column handling and rendering. * [autofix.ci] apply automated fixes * feat: Add unit tests for KBIngestionComponent (#9246) * [autofix.ci] apply automated fixes * fix: remove unnecessary drawer open state change in KnowledgePage * [autofix.ci] apply automated fixes * [autofix.ci] apply automated fixes (attempt 2/3) * Remove kb_info output from KBIngestionComponent (#9275) * [autofix.ci] apply automated fixes * Update Knowledge Bases.json * Use settings service for knowledge base directory Replaces the hardcoded knowledge base directory path with a value from the settings service. This improves configurability and centralizes directory management. * Fix knowledge bases mypy issue * test: Update file page tests for consistency and clarity - Changed expected title text from "My Files" to "Files" for accuracy. - Removed unnecessary parentheses in arrow functions for cleaner syntax. - Updated test assertions to ensure visibility checks are clear and consistent. - Improved readability by standardizing the formatting of test cases. * test: Update expected title in file upload component test for accuracy - Changed expected title text from "My Files" to "Files" to reflect the correct page title. * [autofix.ci] apply automated fixes * Fix tests on backend * Update kb_ingest.py * [autofix.ci] apply automated fixes * Switch to two templates for KB * Update names and descs * [autofix.ci] apply automated fixes * Rename templates * [autofix.ci] apply automated fixes --------- Co-authored-by: Deon Sanchez <69873175+deon-sanchez@users.noreply.github.com> Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com> Co-authored-by: Edwin Jose <edwin.jose@datastax.com>
This commit is contained in:
parent
faff2015c4
commit
e68f6a405a
53 changed files with 7475 additions and 607 deletions
|
|
@ -8,6 +8,7 @@ from langflow.api.v1 import (
|
|||
files_router,
|
||||
flows_router,
|
||||
folders_router,
|
||||
knowledge_bases_router,
|
||||
login_router,
|
||||
mcp_projects_router,
|
||||
mcp_router,
|
||||
|
|
@ -45,6 +46,7 @@ router_v1.include_router(monitor_router)
|
|||
router_v1.include_router(folders_router)
|
||||
router_v1.include_router(projects_router)
|
||||
router_v1.include_router(starter_projects_router)
|
||||
router_v1.include_router(knowledge_bases_router)
|
||||
router_v1.include_router(mcp_router)
|
||||
router_v1.include_router(voice_mode_router)
|
||||
router_v1.include_router(mcp_projects_router)
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@ from langflow.api.v1.endpoints import router as endpoints_router
|
|||
from langflow.api.v1.files import router as files_router
|
||||
from langflow.api.v1.flows import router as flows_router
|
||||
from langflow.api.v1.folders import router as folders_router
|
||||
from langflow.api.v1.knowledge_bases import router as knowledge_bases_router
|
||||
from langflow.api.v1.login import router as login_router
|
||||
from langflow.api.v1.mcp import router as mcp_router
|
||||
from langflow.api.v1.mcp_projects import router as mcp_projects_router
|
||||
|
|
@ -23,6 +24,7 @@ __all__ = [
|
|||
"files_router",
|
||||
"flows_router",
|
||||
"folders_router",
|
||||
"knowledge_bases_router",
|
||||
"login_router",
|
||||
"mcp_projects_router",
|
||||
"mcp_router",
|
||||
|
|
|
|||
437
src/backend/base/langflow/api/v1/knowledge_bases.py
Normal file
437
src/backend/base/langflow/api/v1/knowledge_bases.py
Normal file
|
|
@ -0,0 +1,437 @@
|
|||
import json
|
||||
import shutil
|
||||
from http import HTTPStatus
|
||||
from pathlib import Path
|
||||
|
||||
import pandas as pd
|
||||
from fastapi import APIRouter, HTTPException
|
||||
from langchain_chroma import Chroma
|
||||
from loguru import logger
|
||||
from pydantic import BaseModel
|
||||
|
||||
from langflow.services.deps import get_settings_service
|
||||
|
||||
router = APIRouter(tags=["Knowledge Bases"], prefix="/knowledge_bases")
|
||||
|
||||
|
||||
settings = get_settings_service().settings
|
||||
knowledge_directory = settings.knowledge_bases_dir
|
||||
if not knowledge_directory:
|
||||
msg = "Knowledge bases directory is not set in the settings."
|
||||
raise ValueError(msg)
|
||||
KNOWLEDGE_BASES_DIR = Path(knowledge_directory).expanduser()
|
||||
|
||||
|
||||
class KnowledgeBaseInfo(BaseModel):
|
||||
id: str
|
||||
name: str
|
||||
embedding_provider: str | None = "Unknown"
|
||||
embedding_model: str | None = "Unknown"
|
||||
size: int = 0
|
||||
words: int = 0
|
||||
characters: int = 0
|
||||
chunks: int = 0
|
||||
avg_chunk_size: float = 0.0
|
||||
|
||||
|
||||
class BulkDeleteRequest(BaseModel):
|
||||
kb_names: list[str]
|
||||
|
||||
|
||||
def get_kb_root_path() -> Path:
|
||||
"""Get the knowledge bases root path."""
|
||||
return KNOWLEDGE_BASES_DIR
|
||||
|
||||
|
||||
def get_directory_size(path: Path) -> int:
|
||||
"""Calculate the total size of all files in a directory."""
|
||||
total_size = 0
|
||||
try:
|
||||
for file_path in path.rglob("*"):
|
||||
if file_path.is_file():
|
||||
total_size += file_path.stat().st_size
|
||||
except (OSError, PermissionError):
|
||||
pass
|
||||
return total_size
|
||||
|
||||
|
||||
def detect_embedding_provider(kb_path: Path) -> str:
|
||||
"""Detect the embedding provider from config files and directory structure."""
|
||||
# Provider patterns to check for
|
||||
provider_patterns = {
|
||||
"OpenAI": ["openai", "text-embedding-ada", "text-embedding-3"],
|
||||
"HuggingFace": ["sentence-transformers", "huggingface", "bert-"],
|
||||
"Cohere": ["cohere", "embed-english", "embed-multilingual"],
|
||||
"Google": ["palm", "gecko", "google"],
|
||||
"Chroma": ["chroma"],
|
||||
}
|
||||
|
||||
# Check JSON config files for provider information
|
||||
for config_file in kb_path.glob("*.json"):
|
||||
try:
|
||||
with config_file.open("r", encoding="utf-8") as f:
|
||||
config_data = json.load(f)
|
||||
if not isinstance(config_data, dict):
|
||||
continue
|
||||
|
||||
config_str = json.dumps(config_data).lower()
|
||||
|
||||
# Check for explicit provider fields first
|
||||
provider_fields = ["embedding_provider", "provider", "embedding_model_provider"]
|
||||
for field in provider_fields:
|
||||
if field in config_data:
|
||||
provider_value = str(config_data[field]).lower()
|
||||
for provider, patterns in provider_patterns.items():
|
||||
if any(pattern in provider_value for pattern in patterns):
|
||||
return provider
|
||||
|
||||
# Check for model name patterns
|
||||
for provider, patterns in provider_patterns.items():
|
||||
if any(pattern in config_str for pattern in patterns):
|
||||
return provider
|
||||
|
||||
except (OSError, json.JSONDecodeError) as _:
|
||||
logger.exception("Error reading config file '%s'", config_file)
|
||||
continue
|
||||
|
||||
# Fallback to directory structure
|
||||
if (kb_path / "chroma").exists():
|
||||
return "Chroma"
|
||||
if (kb_path / "vectors.npy").exists():
|
||||
return "Local"
|
||||
|
||||
return "Unknown"
|
||||
|
||||
|
||||
def detect_embedding_model(kb_path: Path) -> str:
|
||||
"""Detect the embedding model from config files."""
|
||||
# First check the embedding metadata file (most accurate)
|
||||
metadata_file = kb_path / "embedding_metadata.json"
|
||||
if metadata_file.exists():
|
||||
try:
|
||||
with metadata_file.open("r", encoding="utf-8") as f:
|
||||
metadata = json.load(f)
|
||||
if isinstance(metadata, dict) and "embedding_model" in metadata:
|
||||
# Check for embedding model field
|
||||
model_value = str(metadata.get("embedding_model", "unknown"))
|
||||
if model_value and model_value.lower() != "unknown":
|
||||
return model_value
|
||||
except (OSError, json.JSONDecodeError) as _:
|
||||
logger.exception("Error reading embedding metadata file '%s'", metadata_file)
|
||||
|
||||
# Check other JSON config files for model information
|
||||
for config_file in kb_path.glob("*.json"):
|
||||
# Skip the embedding metadata file since we already checked it
|
||||
if config_file.name == "embedding_metadata.json":
|
||||
continue
|
||||
|
||||
try:
|
||||
with config_file.open("r", encoding="utf-8") as f:
|
||||
config_data = json.load(f)
|
||||
if not isinstance(config_data, dict):
|
||||
continue
|
||||
|
||||
# Check for explicit model fields first and return the actual model name
|
||||
model_fields = ["embedding_model", "model", "embedding_model_name", "model_name"]
|
||||
for field in model_fields:
|
||||
if field in config_data:
|
||||
model_value = str(config_data[field])
|
||||
if model_value and model_value.lower() != "unknown":
|
||||
return model_value
|
||||
|
||||
# Check for OpenAI specific model names
|
||||
if "openai" in json.dumps(config_data).lower():
|
||||
openai_models = ["text-embedding-ada-002", "text-embedding-3-small", "text-embedding-3-large"]
|
||||
config_str = json.dumps(config_data).lower()
|
||||
for model in openai_models:
|
||||
if model in config_str:
|
||||
return model
|
||||
|
||||
# Check for HuggingFace model names (usually in model field)
|
||||
if "model" in config_data:
|
||||
model_name = str(config_data["model"])
|
||||
# Common HuggingFace embedding models
|
||||
hf_patterns = ["sentence-transformers", "all-MiniLM", "all-mpnet", "multi-qa"]
|
||||
if any(pattern in model_name for pattern in hf_patterns):
|
||||
return model_name
|
||||
|
||||
except (OSError, json.JSONDecodeError) as _:
|
||||
logger.exception("Error reading config file '%s'", config_file)
|
||||
continue
|
||||
|
||||
return "Unknown"
|
||||
|
||||
|
||||
def get_text_columns(df: pd.DataFrame, schema_data: list | None = None) -> list[str]:
|
||||
"""Get the text columns to analyze for word/character counts."""
|
||||
# First try schema-defined text columns
|
||||
if schema_data:
|
||||
text_columns = [
|
||||
col["column_name"]
|
||||
for col in schema_data
|
||||
if col.get("vectorize", False) and col.get("data_type") == "string"
|
||||
]
|
||||
if text_columns:
|
||||
return [col for col in text_columns if col in df.columns]
|
||||
|
||||
# Fallback to common text column names
|
||||
common_names = ["text", "content", "document", "chunk"]
|
||||
text_columns = [col for col in df.columns if col.lower() in common_names]
|
||||
if text_columns:
|
||||
return text_columns
|
||||
|
||||
# Last resort: all string columns
|
||||
return [col for col in df.columns if df[col].dtype == "object"]
|
||||
|
||||
|
||||
def calculate_text_metrics(df: pd.DataFrame, text_columns: list[str]) -> tuple[int, int]:
|
||||
"""Calculate total words and characters from text columns."""
|
||||
total_words = 0
|
||||
total_characters = 0
|
||||
|
||||
for col in text_columns:
|
||||
if col not in df.columns:
|
||||
continue
|
||||
|
||||
text_series = df[col].astype(str).fillna("")
|
||||
total_characters += text_series.str.len().sum()
|
||||
total_words += text_series.str.split().str.len().sum()
|
||||
|
||||
return int(total_words), int(total_characters)
|
||||
|
||||
|
||||
def get_kb_metadata(kb_path: Path) -> dict:
|
||||
"""Extract metadata from a knowledge base directory."""
|
||||
metadata: dict[str, float | int | str] = {
|
||||
"chunks": 0,
|
||||
"words": 0,
|
||||
"characters": 0,
|
||||
"avg_chunk_size": 0.0,
|
||||
"embedding_provider": "Unknown",
|
||||
"embedding_model": "Unknown",
|
||||
}
|
||||
|
||||
try:
|
||||
# First check embedding metadata file for accurate provider and model info
|
||||
metadata_file = kb_path / "embedding_metadata.json"
|
||||
if metadata_file.exists():
|
||||
try:
|
||||
with metadata_file.open("r", encoding="utf-8") as f:
|
||||
embedding_metadata = json.load(f)
|
||||
if isinstance(embedding_metadata, dict):
|
||||
if "embedding_provider" in embedding_metadata:
|
||||
metadata["embedding_provider"] = embedding_metadata["embedding_provider"]
|
||||
if "embedding_model" in embedding_metadata:
|
||||
metadata["embedding_model"] = embedding_metadata["embedding_model"]
|
||||
except (OSError, json.JSONDecodeError) as _:
|
||||
logger.exception("Error reading embedding metadata file '%s'", metadata_file)
|
||||
|
||||
# Fallback to detection if not found in metadata file
|
||||
if metadata["embedding_provider"] == "Unknown":
|
||||
metadata["embedding_provider"] = detect_embedding_provider(kb_path)
|
||||
if metadata["embedding_model"] == "Unknown":
|
||||
metadata["embedding_model"] = detect_embedding_model(kb_path)
|
||||
|
||||
# Read schema for text column information
|
||||
schema_data = None
|
||||
schema_file = kb_path / "schema.json"
|
||||
if schema_file.exists():
|
||||
try:
|
||||
with schema_file.open("r", encoding="utf-8") as f:
|
||||
schema_data = json.load(f)
|
||||
if not isinstance(schema_data, list):
|
||||
schema_data = None
|
||||
except (ValueError, TypeError, OSError) as _:
|
||||
logger.exception("Error reading schema file '%s'", schema_file)
|
||||
|
||||
# Create vector store
|
||||
chroma = Chroma(
|
||||
persist_directory=str(kb_path),
|
||||
collection_name=kb_path.name,
|
||||
)
|
||||
|
||||
# Access the raw collection
|
||||
collection = chroma._collection
|
||||
|
||||
# Fetch all documents and metadata
|
||||
results = collection.get(include=["documents", "metadatas"])
|
||||
|
||||
# Convert to pandas DataFrame
|
||||
source_chunks = pd.DataFrame(
|
||||
{
|
||||
"document": results["documents"],
|
||||
"metadata": results["metadatas"],
|
||||
}
|
||||
)
|
||||
|
||||
# Process the source data for metadata
|
||||
try:
|
||||
metadata["chunks"] = len(source_chunks)
|
||||
|
||||
# Get text columns and calculate metrics
|
||||
text_columns = get_text_columns(source_chunks, schema_data)
|
||||
if text_columns:
|
||||
words, characters = calculate_text_metrics(source_chunks, text_columns)
|
||||
metadata["words"] = words
|
||||
metadata["characters"] = characters
|
||||
|
||||
# Calculate average chunk size
|
||||
if int(metadata["chunks"]) > 0:
|
||||
metadata["avg_chunk_size"] = round(int(characters) / int(metadata["chunks"]), 1)
|
||||
|
||||
except (OSError, ValueError, TypeError) as _:
|
||||
logger.exception("Error processing Chroma DB '%s'", kb_path.name)
|
||||
|
||||
except (OSError, ValueError, TypeError) as _:
|
||||
logger.exception("Error processing knowledge base directory '%s'", kb_path)
|
||||
|
||||
return metadata
|
||||
|
||||
|
||||
@router.get("", status_code=HTTPStatus.OK)
|
||||
@router.get("/", status_code=HTTPStatus.OK)
|
||||
async def list_knowledge_bases() -> list[KnowledgeBaseInfo]:
|
||||
"""List all available knowledge bases."""
|
||||
try:
|
||||
kb_root_path = get_kb_root_path()
|
||||
|
||||
if not kb_root_path.exists():
|
||||
return []
|
||||
|
||||
knowledge_bases = []
|
||||
|
||||
for kb_dir in kb_root_path.iterdir():
|
||||
if not kb_dir.is_dir() or kb_dir.name.startswith("."):
|
||||
continue
|
||||
|
||||
try:
|
||||
# Get size of the directory
|
||||
size = get_directory_size(kb_dir)
|
||||
|
||||
# Get metadata from KB files
|
||||
metadata = get_kb_metadata(kb_dir)
|
||||
|
||||
kb_info = KnowledgeBaseInfo(
|
||||
id=kb_dir.name,
|
||||
name=kb_dir.name.replace("_", " ").replace("-", " ").title(),
|
||||
embedding_provider=metadata["embedding_provider"],
|
||||
embedding_model=metadata["embedding_model"],
|
||||
size=size,
|
||||
words=metadata["words"],
|
||||
characters=metadata["characters"],
|
||||
chunks=metadata["chunks"],
|
||||
avg_chunk_size=metadata["avg_chunk_size"],
|
||||
)
|
||||
|
||||
knowledge_bases.append(kb_info)
|
||||
|
||||
except OSError as _:
|
||||
# Log the exception and skip directories that can't be read
|
||||
logger.exception("Error reading knowledge base directory '%s'", kb_dir)
|
||||
continue
|
||||
|
||||
# Sort by name alphabetically
|
||||
knowledge_bases.sort(key=lambda x: x.name)
|
||||
|
||||
except Exception as e:
|
||||
raise HTTPException(status_code=500, detail=f"Error listing knowledge bases: {e!s}") from e
|
||||
else:
|
||||
return knowledge_bases
|
||||
|
||||
|
||||
@router.get("/{kb_name}", status_code=HTTPStatus.OK)
|
||||
async def get_knowledge_base(kb_name: str) -> KnowledgeBaseInfo:
|
||||
"""Get detailed information about a specific knowledge base."""
|
||||
try:
|
||||
kb_root_path = get_kb_root_path()
|
||||
kb_path = kb_root_path / kb_name
|
||||
|
||||
if not kb_path.exists() or not kb_path.is_dir():
|
||||
raise HTTPException(status_code=404, detail=f"Knowledge base '{kb_name}' not found")
|
||||
|
||||
# Get size of the directory
|
||||
size = get_directory_size(kb_path)
|
||||
|
||||
# Get metadata from KB files
|
||||
metadata = get_kb_metadata(kb_path)
|
||||
|
||||
return KnowledgeBaseInfo(
|
||||
id=kb_name,
|
||||
name=kb_name.replace("_", " ").replace("-", " ").title(),
|
||||
embedding_provider=metadata["embedding_provider"],
|
||||
embedding_model=metadata["embedding_model"],
|
||||
size=size,
|
||||
words=metadata["words"],
|
||||
characters=metadata["characters"],
|
||||
chunks=metadata["chunks"],
|
||||
avg_chunk_size=metadata["avg_chunk_size"],
|
||||
)
|
||||
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception as e:
|
||||
raise HTTPException(status_code=500, detail=f"Error getting knowledge base '{kb_name}': {e!s}") from e
|
||||
|
||||
|
||||
@router.delete("/{kb_name}", status_code=HTTPStatus.OK)
|
||||
async def delete_knowledge_base(kb_name: str) -> dict[str, str]:
|
||||
"""Delete a specific knowledge base."""
|
||||
try:
|
||||
kb_root_path = get_kb_root_path()
|
||||
kb_path = kb_root_path / kb_name
|
||||
|
||||
if not kb_path.exists() or not kb_path.is_dir():
|
||||
raise HTTPException(status_code=404, detail=f"Knowledge base '{kb_name}' not found")
|
||||
|
||||
# Delete the entire knowledge base directory
|
||||
shutil.rmtree(kb_path)
|
||||
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception as e:
|
||||
raise HTTPException(status_code=500, detail=f"Error deleting knowledge base '{kb_name}': {e!s}") from e
|
||||
else:
|
||||
return {"message": f"Knowledge base '{kb_name}' deleted successfully"}
|
||||
|
||||
|
||||
@router.delete("", status_code=HTTPStatus.OK)
|
||||
@router.delete("/", status_code=HTTPStatus.OK)
|
||||
async def delete_knowledge_bases_bulk(request: BulkDeleteRequest) -> dict[str, object]:
|
||||
"""Delete multiple knowledge bases."""
|
||||
try:
|
||||
kb_root_path = get_kb_root_path()
|
||||
deleted_count = 0
|
||||
not_found_kbs = []
|
||||
|
||||
for kb_name in request.kb_names:
|
||||
kb_path = kb_root_path / kb_name
|
||||
|
||||
if not kb_path.exists() or not kb_path.is_dir():
|
||||
not_found_kbs.append(kb_name)
|
||||
continue
|
||||
|
||||
try:
|
||||
# Delete the entire knowledge base directory
|
||||
shutil.rmtree(kb_path)
|
||||
deleted_count += 1
|
||||
except (OSError, PermissionError) as e:
|
||||
logger.exception("Error deleting knowledge base '%s': %s", kb_name, e)
|
||||
# Continue with other deletions even if one fails
|
||||
|
||||
if not_found_kbs and deleted_count == 0:
|
||||
raise HTTPException(status_code=404, detail=f"Knowledge bases not found: {', '.join(not_found_kbs)}")
|
||||
|
||||
result = {
|
||||
"message": f"Successfully deleted {deleted_count} knowledge base(s)",
|
||||
"deleted_count": deleted_count,
|
||||
}
|
||||
|
||||
if not_found_kbs:
|
||||
result["not_found"] = ", ".join(not_found_kbs)
|
||||
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception as e:
|
||||
raise HTTPException(status_code=500, detail=f"Error deleting knowledge bases: {e!s}") from e
|
||||
else:
|
||||
return result
|
||||
104
src/backend/base/langflow/base/data/kb_utils.py
Normal file
104
src/backend/base/langflow/base/data/kb_utils.py
Normal file
|
|
@ -0,0 +1,104 @@
|
|||
import math
|
||||
from collections import Counter
|
||||
|
||||
|
||||
def compute_tfidf(documents: list[str], query_terms: list[str]) -> list[float]:
|
||||
"""Compute TF-IDF scores for query terms across a collection of documents.
|
||||
|
||||
Args:
|
||||
documents: List of document strings
|
||||
query_terms: List of query terms to score
|
||||
|
||||
Returns:
|
||||
List of TF-IDF scores for each document
|
||||
"""
|
||||
# Tokenize documents (simple whitespace splitting)
|
||||
tokenized_docs = [doc.lower().split() for doc in documents]
|
||||
n_docs = len(documents)
|
||||
|
||||
# Calculate document frequency for each term
|
||||
document_frequencies = {}
|
||||
for term in query_terms:
|
||||
document_frequencies[term] = sum(1 for doc in tokenized_docs if term.lower() in doc)
|
||||
|
||||
scores = []
|
||||
|
||||
for doc_tokens in tokenized_docs:
|
||||
doc_score = 0.0
|
||||
doc_length = len(doc_tokens)
|
||||
term_counts = Counter(doc_tokens)
|
||||
|
||||
for term in query_terms:
|
||||
term_lower = term.lower()
|
||||
|
||||
# Term frequency (TF)
|
||||
tf = term_counts[term_lower] / doc_length if doc_length > 0 else 0
|
||||
|
||||
# Inverse document frequency (IDF)
|
||||
idf = math.log(n_docs / document_frequencies[term]) if document_frequencies[term] > 0 else 0
|
||||
|
||||
# TF-IDF score
|
||||
doc_score += tf * idf
|
||||
|
||||
scores.append(doc_score)
|
||||
|
||||
return scores
|
||||
|
||||
|
||||
def compute_bm25(documents: list[str], query_terms: list[str], k1: float = 1.2, b: float = 0.75) -> list[float]:
|
||||
"""Compute BM25 scores for query terms across a collection of documents.
|
||||
|
||||
Args:
|
||||
documents: List of document strings
|
||||
query_terms: List of query terms to score
|
||||
k1: Controls term frequency scaling (default: 1.2)
|
||||
b: Controls document length normalization (default: 0.75)
|
||||
|
||||
Returns:
|
||||
List of BM25 scores for each document
|
||||
"""
|
||||
# Tokenize documents
|
||||
tokenized_docs = [doc.lower().split() for doc in documents]
|
||||
n_docs = len(documents)
|
||||
|
||||
# Calculate average document length
|
||||
avg_doc_length = sum(len(doc) for doc in tokenized_docs) / n_docs if n_docs > 0 else 0
|
||||
|
||||
# Handle edge case where all documents are empty
|
||||
if avg_doc_length == 0:
|
||||
return [0.0] * n_docs
|
||||
|
||||
# Calculate document frequency for each term
|
||||
document_frequencies = {}
|
||||
for term in query_terms:
|
||||
document_frequencies[term] = sum(1 for doc in tokenized_docs if term.lower() in doc)
|
||||
|
||||
scores = []
|
||||
|
||||
for doc_tokens in tokenized_docs:
|
||||
doc_score = 0.0
|
||||
doc_length = len(doc_tokens)
|
||||
term_counts = Counter(doc_tokens)
|
||||
|
||||
for term in query_terms:
|
||||
term_lower = term.lower()
|
||||
|
||||
# Term frequency in document
|
||||
tf = term_counts[term_lower]
|
||||
|
||||
# Inverse document frequency (IDF)
|
||||
# Use standard BM25 IDF formula that ensures non-negative values
|
||||
idf = math.log(n_docs / document_frequencies[term]) if document_frequencies[term] > 0 else 0
|
||||
|
||||
# BM25 score calculation
|
||||
numerator = tf * (k1 + 1)
|
||||
denominator = tf + k1 * (1 - b + b * (doc_length / avg_doc_length))
|
||||
|
||||
# Handle division by zero when tf=0 and k1=0
|
||||
term_score = 0 if denominator == 0 else idf * (numerator / denominator)
|
||||
|
||||
doc_score += term_score
|
||||
|
||||
scores.append(doc_score)
|
||||
|
||||
return scores
|
||||
|
|
@ -3,6 +3,8 @@ from .csv_to_data import CSVToDataComponent
|
|||
from .directory import DirectoryComponent
|
||||
from .file import FileComponent
|
||||
from .json_to_data import JSONToDataComponent
|
||||
from .kb_ingest import KBIngestionComponent
|
||||
from .kb_retrieval import KBRetrievalComponent
|
||||
from .news_search import NewsSearchComponent
|
||||
from .rss import RSSReaderComponent
|
||||
from .sql_executor import SQLComponent
|
||||
|
|
@ -16,6 +18,8 @@ __all__ = [
|
|||
"DirectoryComponent",
|
||||
"FileComponent",
|
||||
"JSONToDataComponent",
|
||||
"KBIngestionComponent",
|
||||
"KBRetrievalComponent",
|
||||
"NewsSearchComponent",
|
||||
"RSSReaderComponent",
|
||||
"SQLComponent",
|
||||
|
|
|
|||
585
src/backend/base/langflow/components/data/kb_ingest.py
Normal file
585
src/backend/base/langflow/components/data/kb_ingest.py
Normal file
|
|
@ -0,0 +1,585 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import re
|
||||
import uuid
|
||||
from dataclasses import asdict, dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import pandas as pd
|
||||
from cryptography.fernet import InvalidToken
|
||||
from langchain_chroma import Chroma
|
||||
from loguru import logger
|
||||
|
||||
from langflow.base.models.openai_constants import OPENAI_EMBEDDING_MODEL_NAMES
|
||||
from langflow.custom import Component
|
||||
from langflow.io import BoolInput, DataFrameInput, DropdownInput, IntInput, Output, SecretStrInput, StrInput, TableInput
|
||||
from langflow.schema.data import Data
|
||||
from langflow.schema.dotdict import dotdict # noqa: TC001
|
||||
from langflow.schema.table import EditMode
|
||||
from langflow.services.auth.utils import decrypt_api_key, encrypt_api_key
|
||||
from langflow.services.deps import get_settings_service
|
||||
|
||||
HUGGINGFACE_MODEL_NAMES = ["sentence-transformers/all-MiniLM-L6-v2", "sentence-transformers/all-mpnet-base-v2"]
|
||||
COHERE_MODEL_NAMES = ["embed-english-v3.0", "embed-multilingual-v3.0"]
|
||||
|
||||
settings = get_settings_service().settings
|
||||
knowledge_directory = settings.knowledge_bases_dir
|
||||
if not knowledge_directory:
|
||||
msg = "Knowledge bases directory is not set in the settings."
|
||||
raise ValueError(msg)
|
||||
KNOWLEDGE_BASES_ROOT_PATH = Path(knowledge_directory).expanduser()
|
||||
|
||||
|
||||
class KBIngestionComponent(Component):
|
||||
"""Create or append to Langflow Knowledge from a DataFrame."""
|
||||
|
||||
# ------ UI metadata ---------------------------------------------------
|
||||
display_name = "Knowledge Ingestion"
|
||||
description = "Create or update knowledge in Langflow."
|
||||
icon = "database"
|
||||
name = "KBIngestion"
|
||||
|
||||
@dataclass
|
||||
class NewKnowledgeBaseInput:
|
||||
functionality: str = "create"
|
||||
fields: dict[str, dict] = field(
|
||||
default_factory=lambda: {
|
||||
"data": {
|
||||
"node": {
|
||||
"name": "create_knowledge_base",
|
||||
"description": "Create new knowledge in Langflow.",
|
||||
"display_name": "Create new knowledge",
|
||||
"field_order": ["01_new_kb_name", "02_embedding_model", "03_api_key"],
|
||||
"template": {
|
||||
"01_new_kb_name": StrInput(
|
||||
name="new_kb_name",
|
||||
display_name="Knowledge Name",
|
||||
info="Name of the new knowledge to create.",
|
||||
required=True,
|
||||
),
|
||||
"02_embedding_model": DropdownInput(
|
||||
name="embedding_model",
|
||||
display_name="Model Name",
|
||||
info="Select the embedding model to use for this knowledge base.",
|
||||
required=True,
|
||||
options=OPENAI_EMBEDDING_MODEL_NAMES + HUGGINGFACE_MODEL_NAMES + COHERE_MODEL_NAMES,
|
||||
options_metadata=[{"icon": "OpenAI"} for _ in OPENAI_EMBEDDING_MODEL_NAMES]
|
||||
+ [{"icon": "HuggingFace"} for _ in HUGGINGFACE_MODEL_NAMES]
|
||||
+ [{"icon": "Cohere"} for _ in COHERE_MODEL_NAMES],
|
||||
),
|
||||
"03_api_key": SecretStrInput(
|
||||
name="api_key",
|
||||
display_name="API Key",
|
||||
info="Provider API key for embedding model",
|
||||
required=True,
|
||||
load_from_db=True,
|
||||
),
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
)
|
||||
|
||||
# ------ Inputs --------------------------------------------------------
|
||||
inputs = [
|
||||
DropdownInput(
|
||||
name="knowledge_base",
|
||||
display_name="Knowledge",
|
||||
info="Select the knowledge to load data from.",
|
||||
required=True,
|
||||
options=[
|
||||
str(d.name) for d in KNOWLEDGE_BASES_ROOT_PATH.iterdir() if not d.name.startswith(".") and d.is_dir()
|
||||
]
|
||||
if KNOWLEDGE_BASES_ROOT_PATH.exists()
|
||||
else [],
|
||||
refresh_button=True,
|
||||
dialog_inputs=asdict(NewKnowledgeBaseInput()),
|
||||
),
|
||||
DataFrameInput(
|
||||
name="input_df",
|
||||
display_name="Data",
|
||||
info="Table with all original columns (already chunked / processed).",
|
||||
required=True,
|
||||
),
|
||||
TableInput(
|
||||
name="column_config",
|
||||
display_name="Column Configuration",
|
||||
info="Configure column behavior for the knowledge base.",
|
||||
required=True,
|
||||
table_schema=[
|
||||
{
|
||||
"name": "column_name",
|
||||
"display_name": "Column Name",
|
||||
"type": "str",
|
||||
"description": "Name of the column in the source DataFrame",
|
||||
"edit_mode": EditMode.INLINE,
|
||||
},
|
||||
{
|
||||
"name": "vectorize",
|
||||
"display_name": "Vectorize",
|
||||
"type": "boolean",
|
||||
"description": "Create embeddings for this column",
|
||||
"default": False,
|
||||
"edit_mode": EditMode.INLINE,
|
||||
},
|
||||
{
|
||||
"name": "identifier",
|
||||
"display_name": "Identifier",
|
||||
"type": "boolean",
|
||||
"description": "Use this column as unique identifier",
|
||||
"default": False,
|
||||
"edit_mode": EditMode.INLINE,
|
||||
},
|
||||
],
|
||||
value=[
|
||||
{
|
||||
"column_name": "text",
|
||||
"vectorize": True,
|
||||
"identifier": False,
|
||||
}
|
||||
],
|
||||
),
|
||||
IntInput(
|
||||
name="chunk_size",
|
||||
display_name="Chunk Size",
|
||||
info="Batch size for processing embeddings",
|
||||
advanced=True,
|
||||
value=1000,
|
||||
),
|
||||
SecretStrInput(
|
||||
name="api_key",
|
||||
display_name="Embedding Provider API Key",
|
||||
info="API key for the embedding provider to generate embeddings.",
|
||||
advanced=True,
|
||||
required=False,
|
||||
),
|
||||
BoolInput(
|
||||
name="allow_duplicates",
|
||||
display_name="Allow Duplicates",
|
||||
info="Allow duplicate rows in the knowledge base",
|
||||
advanced=True,
|
||||
value=False,
|
||||
),
|
||||
]
|
||||
|
||||
# ------ Outputs -------------------------------------------------------
|
||||
outputs = [Output(display_name="DataFrame", name="dataframe", method="build_kb_info")]
|
||||
|
||||
# ------ Internal helpers ---------------------------------------------
|
||||
def _get_kb_root(self) -> Path:
|
||||
"""Return the root directory for knowledge bases."""
|
||||
return KNOWLEDGE_BASES_ROOT_PATH
|
||||
|
||||
def _validate_column_config(self, df_source: pd.DataFrame) -> list[dict[str, Any]]:
|
||||
"""Validate column configuration using Structured Output patterns."""
|
||||
if not self.column_config:
|
||||
msg = "Column configuration cannot be empty"
|
||||
raise ValueError(msg)
|
||||
|
||||
# Convert table input to list of dicts (similar to Structured Output)
|
||||
config_list = self.column_config if isinstance(self.column_config, list) else []
|
||||
|
||||
# Validate column names exist in DataFrame
|
||||
df_columns = set(df_source.columns)
|
||||
for config in config_list:
|
||||
col_name = config.get("column_name")
|
||||
if col_name not in df_columns and not self.silent_errors:
|
||||
msg = f"Column '{col_name}' not found in DataFrame. Available columns: {sorted(df_columns)}"
|
||||
self.log(f"Warning: {msg}")
|
||||
raise ValueError(msg)
|
||||
|
||||
return config_list
|
||||
|
||||
def _get_embedding_provider(self, embedding_model: str) -> str:
|
||||
"""Get embedding provider by matching model name to lists."""
|
||||
if embedding_model in OPENAI_EMBEDDING_MODEL_NAMES:
|
||||
return "OpenAI"
|
||||
if embedding_model in HUGGINGFACE_MODEL_NAMES:
|
||||
return "HuggingFace"
|
||||
if embedding_model in COHERE_MODEL_NAMES:
|
||||
return "Cohere"
|
||||
return "Custom"
|
||||
|
||||
def _build_embeddings(self, embedding_model: str, api_key: str):
|
||||
"""Build embedding model using provider patterns."""
|
||||
# Get provider by matching model name to lists
|
||||
provider = self._get_embedding_provider(embedding_model)
|
||||
|
||||
# Validate provider and model
|
||||
if provider == "OpenAI":
|
||||
from langchain_openai import OpenAIEmbeddings
|
||||
|
||||
if not api_key:
|
||||
msg = "OpenAI API key is required when using OpenAI provider"
|
||||
raise ValueError(msg)
|
||||
return OpenAIEmbeddings(
|
||||
model=embedding_model,
|
||||
api_key=api_key,
|
||||
chunk_size=self.chunk_size,
|
||||
)
|
||||
if provider == "HuggingFace":
|
||||
from langchain_huggingface import HuggingFaceEmbeddings
|
||||
|
||||
return HuggingFaceEmbeddings(
|
||||
model=embedding_model,
|
||||
)
|
||||
if provider == "Cohere":
|
||||
from langchain_cohere import CohereEmbeddings
|
||||
|
||||
if not api_key:
|
||||
msg = "Cohere API key is required when using Cohere provider"
|
||||
raise ValueError(msg)
|
||||
return CohereEmbeddings(
|
||||
model=embedding_model,
|
||||
cohere_api_key=api_key,
|
||||
)
|
||||
if provider == "Custom":
|
||||
# For custom embedding models, we would need additional configuration
|
||||
msg = "Custom embedding models not yet supported"
|
||||
raise NotImplementedError(msg)
|
||||
msg = f"Unknown provider: {provider}"
|
||||
raise ValueError(msg)
|
||||
|
||||
def _build_embedding_metadata(self, embedding_model, api_key) -> dict[str, Any]:
|
||||
"""Build embedding model metadata."""
|
||||
# Get provider by matching model name to lists
|
||||
embedding_provider = self._get_embedding_provider(embedding_model)
|
||||
|
||||
api_key_to_save = None
|
||||
if api_key and hasattr(api_key, "get_secret_value"):
|
||||
api_key_to_save = api_key.get_secret_value()
|
||||
elif isinstance(api_key, str):
|
||||
api_key_to_save = api_key
|
||||
|
||||
encrypted_api_key = None
|
||||
if api_key_to_save:
|
||||
settings_service = get_settings_service()
|
||||
try:
|
||||
encrypted_api_key = encrypt_api_key(api_key_to_save, settings_service=settings_service)
|
||||
except (TypeError, ValueError) as e:
|
||||
self.log(f"Could not encrypt API key: {e}")
|
||||
logger.error(f"Could not encrypt API key: {e}")
|
||||
|
||||
return {
|
||||
"embedding_provider": embedding_provider,
|
||||
"embedding_model": embedding_model,
|
||||
"api_key": encrypted_api_key,
|
||||
"api_key_used": bool(api_key),
|
||||
"chunk_size": self.chunk_size,
|
||||
"created_at": datetime.now(timezone.utc).isoformat(),
|
||||
}
|
||||
|
||||
def _save_embedding_metadata(self, kb_path: Path, embedding_model: str, api_key: str) -> None:
|
||||
"""Save embedding model metadata."""
|
||||
embedding_metadata = self._build_embedding_metadata(embedding_model, api_key)
|
||||
metadata_path = kb_path / "embedding_metadata.json"
|
||||
metadata_path.write_text(json.dumps(embedding_metadata, indent=2))
|
||||
|
||||
def _save_kb_files(
|
||||
self,
|
||||
kb_path: Path,
|
||||
config_list: list[dict[str, Any]],
|
||||
) -> None:
|
||||
"""Save KB files using File Component storage patterns."""
|
||||
try:
|
||||
# Create directory (following File Component patterns)
|
||||
kb_path.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# Save column configuration
|
||||
# Only do this if the file doesn't exist already
|
||||
cfg_path = kb_path / "schema.json"
|
||||
if not cfg_path.exists():
|
||||
cfg_path.write_text(json.dumps(config_list, indent=2))
|
||||
|
||||
except Exception as e:
|
||||
if not self.silent_errors:
|
||||
raise
|
||||
self.log(f"Error saving KB files: {e}")
|
||||
|
||||
def _build_column_metadata(self, config_list: list[dict[str, Any]], df_source: pd.DataFrame) -> dict[str, Any]:
|
||||
"""Build detailed column metadata."""
|
||||
metadata: dict[str, Any] = {
|
||||
"total_columns": len(df_source.columns),
|
||||
"mapped_columns": len(config_list),
|
||||
"unmapped_columns": len(df_source.columns) - len(config_list),
|
||||
"columns": [],
|
||||
"summary": {"vectorized_columns": [], "identifier_columns": []},
|
||||
}
|
||||
|
||||
for config in config_list:
|
||||
col_name = config.get("column_name")
|
||||
vectorize = config.get("vectorize") == "True" or config.get("vectorize") is True
|
||||
identifier = config.get("identifier") == "True" or config.get("identifier") is True
|
||||
|
||||
# Add to columns list
|
||||
metadata["columns"].append(
|
||||
{
|
||||
"name": col_name,
|
||||
"vectorize": vectorize,
|
||||
"identifier": identifier,
|
||||
}
|
||||
)
|
||||
|
||||
# Update summary
|
||||
if vectorize:
|
||||
metadata["summary"]["vectorized_columns"].append(col_name)
|
||||
if identifier:
|
||||
metadata["summary"]["identifier_columns"].append(col_name)
|
||||
|
||||
return metadata
|
||||
|
||||
def _create_vector_store(
|
||||
self, df_source: pd.DataFrame, config_list: list[dict[str, Any]], embedding_model: str, api_key: str
|
||||
) -> None:
|
||||
"""Create vector store following Local DB component pattern."""
|
||||
try:
|
||||
# Set up vector store directory
|
||||
base_dir = self._get_kb_root()
|
||||
|
||||
vector_store_dir = base_dir / self.knowledge_base
|
||||
vector_store_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# Create embeddings model
|
||||
embedding_function = self._build_embeddings(embedding_model, api_key)
|
||||
|
||||
# Convert DataFrame to Data objects (following Local DB pattern)
|
||||
data_objects = self._convert_df_to_data_objects(df_source, config_list)
|
||||
|
||||
# Create vector store
|
||||
chroma = Chroma(
|
||||
persist_directory=str(vector_store_dir),
|
||||
embedding_function=embedding_function,
|
||||
collection_name=self.knowledge_base,
|
||||
)
|
||||
|
||||
# Convert Data objects to LangChain Documents
|
||||
documents = []
|
||||
for data_obj in data_objects:
|
||||
doc = data_obj.to_lc_document()
|
||||
documents.append(doc)
|
||||
|
||||
# Add documents to vector store
|
||||
if documents:
|
||||
chroma.add_documents(documents)
|
||||
self.log(f"Added {len(documents)} documents to vector store '{self.knowledge_base}'")
|
||||
|
||||
except Exception as e:
|
||||
if not self.silent_errors:
|
||||
raise
|
||||
self.log(f"Error creating vector store: {e}")
|
||||
|
||||
def _convert_df_to_data_objects(self, df_source: pd.DataFrame, config_list: list[dict[str, Any]]) -> list[Data]:
|
||||
"""Convert DataFrame to Data objects for vector store."""
|
||||
data_objects: list[Data] = []
|
||||
|
||||
# Set up vector store directory
|
||||
base_dir = self._get_kb_root()
|
||||
|
||||
# If we don't allow duplicates, we need to get the existing hashes
|
||||
chroma = Chroma(
|
||||
persist_directory=str(base_dir / self.knowledge_base),
|
||||
collection_name=self.knowledge_base,
|
||||
)
|
||||
|
||||
# Get all documents and their metadata
|
||||
all_docs = chroma.get()
|
||||
|
||||
# Extract all _id values from metadata
|
||||
id_list = [metadata.get("_id") for metadata in all_docs["metadatas"] if metadata.get("_id")]
|
||||
|
||||
# Get column roles
|
||||
content_cols = []
|
||||
identifier_cols = []
|
||||
|
||||
for config in config_list:
|
||||
col_name = config.get("column_name")
|
||||
vectorize = config.get("vectorize") == "True" or config.get("vectorize") is True
|
||||
identifier = config.get("identifier") == "True" or config.get("identifier") is True
|
||||
|
||||
if vectorize:
|
||||
content_cols.append(col_name)
|
||||
elif identifier:
|
||||
identifier_cols.append(col_name)
|
||||
|
||||
# Convert each row to a Data object
|
||||
for _, row in df_source.iterrows():
|
||||
# Build content text from vectorized columns using list comprehension
|
||||
content_parts = [str(row[col]) for col in content_cols if col in row and pd.notna(row[col])]
|
||||
|
||||
page_content = " ".join(content_parts)
|
||||
|
||||
# Build metadata from NON-vectorized columns only (simple key-value pairs)
|
||||
data_dict = {
|
||||
"text": page_content, # Main content for vectorization
|
||||
}
|
||||
|
||||
# Add metadata columns as simple key-value pairs
|
||||
for col in df_source.columns:
|
||||
if col not in content_cols and col in row and pd.notna(row[col]):
|
||||
# Convert to simple types for Chroma metadata
|
||||
value = row[col]
|
||||
data_dict[col] = str(value) # Convert complex types to string
|
||||
|
||||
# Hash the page_content for unique ID
|
||||
page_content_hash = hashlib.sha256(page_content.encode()).hexdigest()
|
||||
data_dict["_id"] = page_content_hash
|
||||
|
||||
# If duplicates are disallowed, and hash exists, prevent adding this row
|
||||
if not self.allow_duplicates and page_content_hash in id_list:
|
||||
self.log(f"Skipping duplicate row with hash {page_content_hash}")
|
||||
continue
|
||||
|
||||
# Create Data object - everything except "text" becomes metadata
|
||||
data_obj = Data(data=data_dict)
|
||||
data_objects.append(data_obj)
|
||||
|
||||
return data_objects
|
||||
|
||||
def is_valid_collection_name(self, name, min_length: int = 3, max_length: int = 63) -> bool:
|
||||
"""Validates collection name against conditions 1-3.
|
||||
|
||||
1. Contains 3-63 characters
|
||||
2. Starts and ends with alphanumeric character
|
||||
3. Contains only alphanumeric characters, underscores, or hyphens.
|
||||
|
||||
Args:
|
||||
name (str): Collection name to validate
|
||||
min_length (int): Minimum length of the name
|
||||
max_length (int): Maximum length of the name
|
||||
|
||||
Returns:
|
||||
bool: True if valid, False otherwise
|
||||
"""
|
||||
# Check length (condition 1)
|
||||
if not (min_length <= len(name) <= max_length):
|
||||
return False
|
||||
|
||||
# Check start/end with alphanumeric (condition 2)
|
||||
if not (name[0].isalnum() and name[-1].isalnum()):
|
||||
return False
|
||||
|
||||
# Check allowed characters (condition 3)
|
||||
return re.match(r"^[a-zA-Z0-9_-]+$", name) is not None
|
||||
|
||||
# ---------------------------------------------------------------------
|
||||
# OUTPUT METHODS
|
||||
# ---------------------------------------------------------------------
|
||||
def build_kb_info(self) -> Data:
|
||||
"""Main ingestion routine → returns a dict with KB metadata."""
|
||||
try:
|
||||
# Get source DataFrame
|
||||
df_source: pd.DataFrame = self.input_df
|
||||
|
||||
# Validate column configuration (using Structured Output patterns)
|
||||
config_list = self._validate_column_config(df_source)
|
||||
column_metadata = self._build_column_metadata(config_list, df_source)
|
||||
|
||||
# Prepare KB folder (using File Component patterns)
|
||||
kb_root = self._get_kb_root()
|
||||
kb_path = kb_root / self.knowledge_base
|
||||
|
||||
# Read the embedding info from the knowledge base folder
|
||||
metadata_path = kb_path / "embedding_metadata.json"
|
||||
|
||||
# If the API key is not provided, try to read it from the metadata file
|
||||
if metadata_path.exists():
|
||||
settings_service = get_settings_service()
|
||||
metadata = json.loads(metadata_path.read_text())
|
||||
embedding_model = metadata.get("embedding_model")
|
||||
try:
|
||||
api_key = decrypt_api_key(metadata["api_key"], settings_service)
|
||||
except (InvalidToken, TypeError, ValueError) as e:
|
||||
logger.error(f"Could not decrypt API key. Please provide it manually. Error: {e}")
|
||||
|
||||
# Check if a custom API key was provided, update metadata if so
|
||||
if self.api_key:
|
||||
api_key = self.api_key
|
||||
self._save_embedding_metadata(
|
||||
kb_path=kb_path,
|
||||
embedding_model=embedding_model,
|
||||
api_key=api_key,
|
||||
)
|
||||
|
||||
# Create vector store following Local DB component pattern
|
||||
self._create_vector_store(df_source, config_list, embedding_model=embedding_model, api_key=api_key)
|
||||
|
||||
# Save KB files (using File Component storage patterns)
|
||||
self._save_kb_files(kb_path, config_list)
|
||||
|
||||
# Build metadata response
|
||||
meta: dict[str, Any] = {
|
||||
"kb_id": str(uuid.uuid4()),
|
||||
"kb_name": self.knowledge_base,
|
||||
"rows": len(df_source),
|
||||
"column_metadata": column_metadata,
|
||||
"path": str(kb_path),
|
||||
"config_columns": len(config_list),
|
||||
"timestamp": datetime.now(tz=timezone.utc).isoformat(),
|
||||
}
|
||||
|
||||
# Set status message
|
||||
self.status = f"✅ KB **{self.knowledge_base}** saved · {len(df_source)} chunks."
|
||||
|
||||
return Data(data=meta)
|
||||
|
||||
except Exception as e:
|
||||
if not self.silent_errors:
|
||||
raise
|
||||
self.log(f"Error in KB ingestion: {e}")
|
||||
self.status = f"❌ KB ingestion failed: {e}"
|
||||
return Data(data={"error": str(e), "kb_name": self.knowledge_base})
|
||||
|
||||
def _get_knowledge_bases(self) -> list[str]:
|
||||
"""Retrieve a list of available knowledge bases.
|
||||
|
||||
Returns:
|
||||
A list of knowledge base names.
|
||||
"""
|
||||
# Return the list of directories in the knowledge base root path
|
||||
kb_root_path = self._get_kb_root()
|
||||
|
||||
if not kb_root_path.exists():
|
||||
return []
|
||||
|
||||
return [str(d.name) for d in kb_root_path.iterdir() if not d.name.startswith(".") and d.is_dir()]
|
||||
|
||||
def update_build_config(self, build_config: dotdict, field_value: Any, field_name: str | None = None) -> dotdict:
|
||||
"""Update build configuration based on provider selection."""
|
||||
# Create a new knowledge base
|
||||
if field_name == "knowledge_base":
|
||||
if isinstance(field_value, dict) and "01_new_kb_name" in field_value:
|
||||
# Validate the knowledge base name - Make sure it follows these rules:
|
||||
if not self.is_valid_collection_name(field_value["01_new_kb_name"]):
|
||||
msg = f"Invalid knowledge base name: {field_value['01_new_kb_name']}"
|
||||
raise ValueError(msg)
|
||||
|
||||
# We need to test the API Key one time against the embedding model
|
||||
embed_model = self._build_embeddings(
|
||||
embedding_model=field_value["02_embedding_model"], api_key=field_value["03_api_key"]
|
||||
)
|
||||
|
||||
# Try to generate a dummy embedding to validate the API key
|
||||
embed_model.embed_query("test")
|
||||
|
||||
# Create the new knowledge base directory
|
||||
kb_path = KNOWLEDGE_BASES_ROOT_PATH / field_value["01_new_kb_name"]
|
||||
kb_path.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# Save the embedding metadata
|
||||
build_config["knowledge_base"]["value"] = field_value["01_new_kb_name"]
|
||||
self._save_embedding_metadata(
|
||||
kb_path=kb_path,
|
||||
embedding_model=field_value["02_embedding_model"],
|
||||
api_key=field_value["03_api_key"],
|
||||
)
|
||||
|
||||
# Update the knowledge base options dynamically
|
||||
build_config["knowledge_base"]["options"] = self._get_knowledge_bases()
|
||||
if build_config["knowledge_base"]["value"] not in build_config["knowledge_base"]["options"]:
|
||||
build_config["knowledge_base"]["value"] = None
|
||||
|
||||
return build_config
|
||||
254
src/backend/base/langflow/components/data/kb_retrieval.py
Normal file
254
src/backend/base/langflow/components/data/kb_retrieval.py
Normal file
|
|
@ -0,0 +1,254 @@
|
|||
import json
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from cryptography.fernet import InvalidToken
|
||||
from langchain_chroma import Chroma
|
||||
from loguru import logger
|
||||
|
||||
from langflow.custom import Component
|
||||
from langflow.io import BoolInput, DropdownInput, IntInput, MessageTextInput, Output, SecretStrInput
|
||||
from langflow.schema.data import Data
|
||||
from langflow.schema.dataframe import DataFrame
|
||||
from langflow.services.auth.utils import decrypt_api_key
|
||||
from langflow.services.deps import get_settings_service
|
||||
|
||||
settings = get_settings_service().settings
|
||||
knowledge_directory = settings.knowledge_bases_dir
|
||||
if not knowledge_directory:
|
||||
msg = "Knowledge bases directory is not set in the settings."
|
||||
raise ValueError(msg)
|
||||
KNOWLEDGE_BASES_ROOT_PATH = Path(knowledge_directory).expanduser()
|
||||
|
||||
|
||||
class KBRetrievalComponent(Component):
|
||||
display_name = "Knowledge Retrieval"
|
||||
description = "Search and retrieve data from knowledge."
|
||||
icon = "database"
|
||||
name = "KBRetrieval"
|
||||
|
||||
inputs = [
|
||||
DropdownInput(
|
||||
name="knowledge_base",
|
||||
display_name="Knowledge",
|
||||
info="Select the knowledge to load data from.",
|
||||
required=True,
|
||||
options=[
|
||||
str(d.name) for d in KNOWLEDGE_BASES_ROOT_PATH.iterdir() if not d.name.startswith(".") and d.is_dir()
|
||||
]
|
||||
if KNOWLEDGE_BASES_ROOT_PATH.exists()
|
||||
else [],
|
||||
refresh_button=True,
|
||||
real_time_refresh=True,
|
||||
),
|
||||
SecretStrInput(
|
||||
name="api_key",
|
||||
display_name="Embedding Provider API Key",
|
||||
info="API key for the embedding provider to generate embeddings.",
|
||||
advanced=True,
|
||||
required=False,
|
||||
),
|
||||
MessageTextInput(
|
||||
name="search_query",
|
||||
display_name="Search Query",
|
||||
info="Optional search query to filter knowledge base data.",
|
||||
),
|
||||
IntInput(
|
||||
name="top_k",
|
||||
display_name="Top K Results",
|
||||
info="Number of top results to return from the knowledge base.",
|
||||
value=5,
|
||||
advanced=True,
|
||||
required=False,
|
||||
),
|
||||
BoolInput(
|
||||
name="include_metadata",
|
||||
display_name="Include Metadata",
|
||||
info="Whether to include all metadata and embeddings in the output. If false, only content is returned.",
|
||||
value=True,
|
||||
advanced=True,
|
||||
),
|
||||
]
|
||||
|
||||
outputs = [
|
||||
Output(
|
||||
name="chroma_kb_data",
|
||||
display_name="Results",
|
||||
method="get_chroma_kb_data",
|
||||
info="Returns the data from the selected knowledge base.",
|
||||
),
|
||||
]
|
||||
|
||||
def _get_knowledge_bases(self) -> list[str]:
|
||||
"""Retrieve a list of available knowledge bases.
|
||||
|
||||
Returns:
|
||||
A list of knowledge base names.
|
||||
"""
|
||||
if not KNOWLEDGE_BASES_ROOT_PATH.exists():
|
||||
return []
|
||||
|
||||
return [str(d.name) for d in KNOWLEDGE_BASES_ROOT_PATH.iterdir() if not d.name.startswith(".") and d.is_dir()]
|
||||
|
||||
def update_build_config(self, build_config, field_value, field_name=None): # noqa: ARG002
|
||||
if field_name == "knowledge_base":
|
||||
# Update the knowledge base options dynamically
|
||||
build_config["knowledge_base"]["options"] = self._get_knowledge_bases()
|
||||
|
||||
# If the selected knowledge base is not available, reset it
|
||||
if build_config["knowledge_base"]["value"] not in build_config["knowledge_base"]["options"]:
|
||||
build_config["knowledge_base"]["value"] = None
|
||||
|
||||
return build_config
|
||||
|
||||
def _get_kb_metadata(self, kb_path: Path) -> dict:
|
||||
"""Load and process knowledge base metadata."""
|
||||
metadata: dict[str, Any] = {}
|
||||
metadata_file = kb_path / "embedding_metadata.json"
|
||||
if not metadata_file.exists():
|
||||
logger.warning(f"Embedding metadata file not found at {metadata_file}")
|
||||
return metadata
|
||||
|
||||
try:
|
||||
with metadata_file.open("r", encoding="utf-8") as f:
|
||||
metadata = json.load(f)
|
||||
except json.JSONDecodeError:
|
||||
logger.error(f"Error decoding JSON from {metadata_file}")
|
||||
return {}
|
||||
|
||||
# Decrypt API key if it exists
|
||||
if "api_key" in metadata and metadata.get("api_key"):
|
||||
settings_service = get_settings_service()
|
||||
try:
|
||||
decrypted_key = decrypt_api_key(metadata["api_key"], settings_service)
|
||||
metadata["api_key"] = decrypted_key
|
||||
except (InvalidToken, TypeError, ValueError) as e:
|
||||
logger.error(f"Could not decrypt API key. Please provide it manually. Error: {e}")
|
||||
metadata["api_key"] = None
|
||||
return metadata
|
||||
|
||||
def _build_embeddings(self, metadata: dict):
|
||||
"""Build embedding model from metadata."""
|
||||
provider = metadata.get("embedding_provider")
|
||||
model = metadata.get("embedding_model")
|
||||
api_key = metadata.get("api_key")
|
||||
chunk_size = metadata.get("chunk_size")
|
||||
|
||||
# If user provided a key in the input, it overrides the stored one.
|
||||
if self.api_key and self.api_key.get_secret_value():
|
||||
api_key = self.api_key.get_secret_value()
|
||||
|
||||
# Handle various providers
|
||||
if provider == "OpenAI":
|
||||
from langchain_openai import OpenAIEmbeddings
|
||||
|
||||
if not api_key:
|
||||
msg = "OpenAI API key is required. Provide it in the component's advanced settings."
|
||||
raise ValueError(msg)
|
||||
return OpenAIEmbeddings(
|
||||
model=model,
|
||||
api_key=api_key,
|
||||
chunk_size=chunk_size,
|
||||
)
|
||||
if provider == "HuggingFace":
|
||||
from langchain_huggingface import HuggingFaceEmbeddings
|
||||
|
||||
return HuggingFaceEmbeddings(
|
||||
model=model,
|
||||
)
|
||||
if provider == "Cohere":
|
||||
from langchain_cohere import CohereEmbeddings
|
||||
|
||||
if not api_key:
|
||||
msg = "Cohere API key is required when using Cohere provider"
|
||||
raise ValueError(msg)
|
||||
return CohereEmbeddings(
|
||||
model=model,
|
||||
cohere_api_key=api_key,
|
||||
)
|
||||
if provider == "Custom":
|
||||
# For custom embedding models, we would need additional configuration
|
||||
msg = "Custom embedding models not yet supported"
|
||||
raise NotImplementedError(msg)
|
||||
# Add other providers here if they become supported in ingest
|
||||
msg = f"Embedding provider '{provider}' is not supported for retrieval."
|
||||
raise NotImplementedError(msg)
|
||||
|
||||
def get_chroma_kb_data(self) -> DataFrame:
|
||||
"""Retrieve data from the selected knowledge base by reading the Chroma collection.
|
||||
|
||||
Returns:
|
||||
A DataFrame containing the data rows from the knowledge base.
|
||||
"""
|
||||
kb_path = KNOWLEDGE_BASES_ROOT_PATH / self.knowledge_base
|
||||
|
||||
metadata = self._get_kb_metadata(kb_path)
|
||||
if not metadata:
|
||||
msg = f"Metadata not found for knowledge base: {self.knowledge_base}. Ensure it has been indexed."
|
||||
raise ValueError(msg)
|
||||
|
||||
# Build the embedder for the knowledge base
|
||||
embedding_function = self._build_embeddings(metadata)
|
||||
|
||||
# Load vector store
|
||||
chroma = Chroma(
|
||||
persist_directory=str(kb_path),
|
||||
embedding_function=embedding_function,
|
||||
collection_name=self.knowledge_base,
|
||||
)
|
||||
|
||||
# If a search query is provided, perform a similarity search
|
||||
if self.search_query:
|
||||
# Use the search query to perform a similarity search
|
||||
logger.info(f"Performing similarity search with query: {self.search_query}")
|
||||
results = chroma.similarity_search_with_score(
|
||||
query=self.search_query or "",
|
||||
k=self.top_k,
|
||||
)
|
||||
else:
|
||||
results = chroma.similarity_search(
|
||||
query=self.search_query or "",
|
||||
k=self.top_k,
|
||||
)
|
||||
|
||||
# For each result, make it a tuple to match the expected output format
|
||||
results = [(doc, 0) for doc in results] # Assign a dummy score of 0
|
||||
|
||||
# If metadata is enabled, get embeddings for the results
|
||||
id_to_embedding = {}
|
||||
if self.include_metadata and results:
|
||||
doc_ids = [doc[0].metadata.get("_id") for doc in results if doc[0].metadata.get("_id")]
|
||||
|
||||
# Only proceed if we have valid document IDs
|
||||
if doc_ids:
|
||||
# Access underlying client to get embeddings
|
||||
collection = chroma._client.get_collection(name=self.knowledge_base)
|
||||
embeddings_result = collection.get(where={"_id": {"$in": doc_ids}}, include=["embeddings", "metadatas"])
|
||||
|
||||
# Create a mapping from document ID to embedding
|
||||
for i, metadata in enumerate(embeddings_result.get("metadatas", [])):
|
||||
if metadata and "_id" in metadata:
|
||||
id_to_embedding[metadata["_id"]] = embeddings_result["embeddings"][i]
|
||||
|
||||
# Build output data based on include_metadata setting
|
||||
data_list = []
|
||||
for doc in results:
|
||||
if self.include_metadata:
|
||||
# Include all metadata, embeddings, and content
|
||||
kwargs = {
|
||||
"content": doc[0].page_content,
|
||||
**doc[0].metadata,
|
||||
}
|
||||
if self.search_query:
|
||||
kwargs["_score"] = -1 * doc[1]
|
||||
kwargs["_embeddings"] = id_to_embedding.get(doc[0].metadata.get("_id"))
|
||||
else:
|
||||
# Only include content
|
||||
kwargs = {
|
||||
"content": doc[0].page_content,
|
||||
}
|
||||
|
||||
data_list.append(Data(**kwargs))
|
||||
|
||||
# Return the DataFrame containing the data
|
||||
return DataFrame(data=data_list)
|
||||
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
|
|
@ -73,6 +73,9 @@ class Settings(BaseSettings):
|
|||
"""Define if langflow database should be saved in LANGFLOW_CONFIG_DIR or in the langflow directory
|
||||
(i.e. in the package directory)."""
|
||||
|
||||
knowledge_bases_dir: str | None = "~/.langflow/knowledge_bases"
|
||||
"""The directory to store knowledge bases."""
|
||||
|
||||
dev: bool = False
|
||||
"""If True, Langflow will run in development mode."""
|
||||
database_url: str | None = None
|
||||
|
|
|
|||
0
src/backend/tests/unit/base/data/__init__.py
Normal file
0
src/backend/tests/unit/base/data/__init__.py
Normal file
458
src/backend/tests/unit/base/data/test_kb_utils.py
Normal file
458
src/backend/tests/unit/base/data/test_kb_utils.py
Normal file
|
|
@ -0,0 +1,458 @@
|
|||
import pytest
|
||||
from langflow.base.data.kb_utils import compute_bm25, compute_tfidf
|
||||
|
||||
|
||||
class TestKBUtils:
|
||||
"""Test suite for knowledge base utility functions."""
|
||||
|
||||
# Test data for TF-IDF and BM25 tests
|
||||
@pytest.fixture
|
||||
def sample_documents(self):
|
||||
"""Sample documents for testing."""
|
||||
return ["the cat sat on the mat", "the dog ran in the park", "cats and dogs are pets", "birds fly in the sky"]
|
||||
|
||||
@pytest.fixture
|
||||
def query_terms(self):
|
||||
"""Sample query terms for testing."""
|
||||
return ["cat", "dog"]
|
||||
|
||||
@pytest.fixture
|
||||
def empty_documents(self):
|
||||
"""Empty documents for edge case testing."""
|
||||
return ["", "", ""]
|
||||
|
||||
@pytest.fixture
|
||||
def single_document(self):
|
||||
"""Single document for testing."""
|
||||
return ["hello world"]
|
||||
|
||||
def test_compute_tfidf_basic(self, sample_documents, query_terms):
|
||||
"""Test basic TF-IDF computation."""
|
||||
scores = compute_tfidf(sample_documents, query_terms)
|
||||
|
||||
# Should return a score for each document
|
||||
assert len(scores) == len(sample_documents)
|
||||
|
||||
# All scores should be floats
|
||||
assert all(isinstance(score, float) for score in scores)
|
||||
|
||||
# First document contains "cat", should have non-zero score
|
||||
assert scores[0] > 0.0
|
||||
|
||||
# Second document contains "dog", should have non-zero score
|
||||
assert scores[1] > 0.0
|
||||
|
||||
# Third document contains both "cats" and "dogs", but case-insensitive matching should work
|
||||
# Note: "cats" != "cat" exactly, so this tests the term matching behavior
|
||||
assert scores[2] >= 0.0
|
||||
|
||||
# Fourth document contains neither term, should have zero score
|
||||
assert scores[3] == 0.0
|
||||
|
||||
def test_compute_tfidf_case_insensitive(self):
|
||||
"""Test that TF-IDF computation is case insensitive."""
|
||||
documents = ["The CAT sat", "the dog RAN", "CATS and DOGS"]
|
||||
query_terms = ["cat", "DOG"]
|
||||
|
||||
scores = compute_tfidf(documents, query_terms)
|
||||
|
||||
# First document should match "cat" (case insensitive)
|
||||
assert scores[0] > 0.0
|
||||
|
||||
# Second document should match "dog" (case insensitive)
|
||||
assert scores[1] > 0.0
|
||||
|
||||
def test_compute_tfidf_empty_documents(self, empty_documents, query_terms):
|
||||
"""Test TF-IDF with empty documents."""
|
||||
scores = compute_tfidf(empty_documents, query_terms)
|
||||
|
||||
# Should return scores for all documents
|
||||
assert len(scores) == len(empty_documents)
|
||||
|
||||
# All scores should be zero since documents are empty
|
||||
assert all(score == 0.0 for score in scores)
|
||||
|
||||
def test_compute_tfidf_empty_query_terms(self, sample_documents):
|
||||
"""Test TF-IDF with empty query terms."""
|
||||
scores = compute_tfidf(sample_documents, [])
|
||||
|
||||
# Should return scores for all documents
|
||||
assert len(scores) == len(sample_documents)
|
||||
|
||||
# All scores should be zero since no query terms
|
||||
assert all(score == 0.0 for score in scores)
|
||||
|
||||
def test_compute_tfidf_single_document(self, single_document):
|
||||
"""Test TF-IDF with single document."""
|
||||
query_terms = ["hello", "world"]
|
||||
scores = compute_tfidf(single_document, query_terms)
|
||||
|
||||
assert len(scores) == 1
|
||||
# With only one document, IDF = log(1/1) = 0, so TF-IDF score is always 0
|
||||
# This is correct mathematical behavior - TF-IDF is designed to discriminate between documents
|
||||
assert scores[0] == 0.0
|
||||
|
||||
def test_compute_tfidf_two_documents_positive_scores(self):
|
||||
"""Test TF-IDF with two documents to ensure positive scores are possible."""
|
||||
documents = ["hello world", "goodbye earth"]
|
||||
query_terms = ["hello", "world"]
|
||||
scores = compute_tfidf(documents, query_terms)
|
||||
|
||||
assert len(scores) == 2
|
||||
# First document contains both terms, should have positive score
|
||||
assert scores[0] > 0.0
|
||||
# Second document contains neither term, should have zero score
|
||||
assert scores[1] == 0.0
|
||||
|
||||
def test_compute_tfidf_no_documents(self):
|
||||
"""Test TF-IDF with no documents."""
|
||||
scores = compute_tfidf([], ["cat", "dog"])
|
||||
|
||||
assert scores == []
|
||||
|
||||
def test_compute_tfidf_term_frequency_calculation(self):
|
||||
"""Test TF-IDF term frequency calculation."""
|
||||
# Documents with different term frequencies for the same term
|
||||
documents = ["rare word text", "rare rare word", "other content"]
|
||||
query_terms = ["rare"]
|
||||
|
||||
scores = compute_tfidf(documents, query_terms)
|
||||
|
||||
# "rare" appears in documents 0 and 1, but with different frequencies
|
||||
# Document 1 has higher TF (2/3 vs 1/3), so should score higher
|
||||
assert scores[0] > 0.0 # Contains "rare" once
|
||||
assert scores[1] > scores[0] # Contains "rare" twice, should score higher
|
||||
assert scores[2] == 0.0 # Doesn't contain "rare"
|
||||
|
||||
def test_compute_tfidf_idf_calculation(self):
|
||||
"""Test TF-IDF inverse document frequency calculation."""
|
||||
# "rare" appears in only one document, "common" appears in both
|
||||
documents = ["rare term", "common term", "common word"]
|
||||
query_terms = ["rare", "common"]
|
||||
|
||||
scores = compute_tfidf(documents, query_terms)
|
||||
|
||||
# First document should have higher score due to rare term having higher IDF
|
||||
assert scores[0] > scores[1] # rare term gets higher IDF
|
||||
assert scores[0] > scores[2]
|
||||
|
||||
def test_compute_bm25_basic(self, sample_documents, query_terms):
|
||||
"""Test basic BM25 computation."""
|
||||
scores = compute_bm25(sample_documents, query_terms)
|
||||
|
||||
# Should return a score for each document
|
||||
assert len(scores) == len(sample_documents)
|
||||
|
||||
# All scores should be floats
|
||||
assert all(isinstance(score, float) for score in scores)
|
||||
|
||||
# First document contains "cat", should have non-zero score
|
||||
assert scores[0] > 0.0
|
||||
|
||||
# Second document contains "dog", should have non-zero score
|
||||
assert scores[1] > 0.0
|
||||
|
||||
# Fourth document contains neither term, should have zero score
|
||||
assert scores[3] == 0.0
|
||||
|
||||
def test_compute_bm25_parameters(self, sample_documents, query_terms):
|
||||
"""Test BM25 with different k1 and b parameters."""
|
||||
# Test with default parameters
|
||||
scores_default = compute_bm25(sample_documents, query_terms)
|
||||
|
||||
# Test with different k1
|
||||
scores_k1 = compute_bm25(sample_documents, query_terms, k1=2.0)
|
||||
|
||||
# Test with different b
|
||||
scores_b = compute_bm25(sample_documents, query_terms, b=0.5)
|
||||
|
||||
# Test with both different
|
||||
scores_both = compute_bm25(sample_documents, query_terms, k1=2.0, b=0.5)
|
||||
|
||||
# All should return valid scores
|
||||
assert len(scores_default) == len(sample_documents)
|
||||
assert len(scores_k1) == len(sample_documents)
|
||||
assert len(scores_b) == len(sample_documents)
|
||||
assert len(scores_both) == len(sample_documents)
|
||||
|
||||
# Scores should be different with different parameters
|
||||
assert scores_default != scores_k1
|
||||
assert scores_default != scores_b
|
||||
|
||||
def test_compute_bm25_case_insensitive(self):
|
||||
"""Test that BM25 computation is case insensitive."""
|
||||
documents = ["The CAT sat", "the dog RAN", "CATS and DOGS"]
|
||||
query_terms = ["cat", "DOG"]
|
||||
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
# First document should match "cat" (case insensitive)
|
||||
assert scores[0] > 0.0
|
||||
|
||||
# Second document should match "dog" (case insensitive)
|
||||
assert scores[1] > 0.0
|
||||
|
||||
def test_compute_bm25_empty_documents(self, empty_documents, query_terms):
|
||||
"""Test BM25 with empty documents."""
|
||||
scores = compute_bm25(empty_documents, query_terms)
|
||||
|
||||
# Should return scores for all documents
|
||||
assert len(scores) == len(empty_documents)
|
||||
|
||||
# All scores should be zero since documents are empty
|
||||
assert all(score == 0.0 for score in scores)
|
||||
|
||||
def test_compute_bm25_empty_query_terms(self, sample_documents):
|
||||
"""Test BM25 with empty query terms."""
|
||||
scores = compute_bm25(sample_documents, [])
|
||||
|
||||
# Should return scores for all documents
|
||||
assert len(scores) == len(sample_documents)
|
||||
|
||||
# All scores should be zero since no query terms
|
||||
assert all(score == 0.0 for score in scores)
|
||||
|
||||
def test_compute_bm25_single_document(self, single_document):
|
||||
"""Test BM25 with single document."""
|
||||
query_terms = ["hello", "world"]
|
||||
scores = compute_bm25(single_document, query_terms)
|
||||
|
||||
assert len(scores) == 1
|
||||
# With only one document, IDF = log(1/1) = 0, so BM25 score is always 0
|
||||
# This is correct mathematical behavior - both TF-IDF and BM25 are designed to discriminate between documents
|
||||
assert scores[0] == 0.0
|
||||
|
||||
def test_compute_bm25_two_documents_positive_scores(self):
|
||||
"""Test BM25 with two documents to ensure positive scores are possible."""
|
||||
documents = ["hello world", "goodbye earth"]
|
||||
query_terms = ["hello", "world"]
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
assert len(scores) == 2
|
||||
# First document contains both terms, should have positive score
|
||||
assert scores[0] > 0.0
|
||||
# Second document contains neither term, should have zero score
|
||||
assert scores[1] == 0.0
|
||||
|
||||
def test_compute_bm25_no_documents(self):
|
||||
"""Test BM25 with no documents."""
|
||||
scores = compute_bm25([], ["cat", "dog"])
|
||||
|
||||
assert scores == []
|
||||
|
||||
def test_compute_bm25_document_length_normalization(self):
|
||||
"""Test BM25 document length normalization."""
|
||||
# Test with documents where some terms appear in subset of documents
|
||||
documents = [
|
||||
"cat unique1", # Short document with unique term
|
||||
"cat dog bird mouse elephant tiger lion bear wolf unique2", # Long document with unique term
|
||||
"other content", # Document without query terms
|
||||
]
|
||||
query_terms = ["unique1", "unique2"]
|
||||
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
# Documents with unique terms should have positive scores
|
||||
assert scores[0] > 0.0 # Contains "unique1"
|
||||
assert scores[1] > 0.0 # Contains "unique2"
|
||||
assert scores[2] == 0.0 # Contains neither term
|
||||
|
||||
# Document length normalization affects scores
|
||||
assert len(scores) == 3
|
||||
|
||||
def test_compute_bm25_term_frequency_saturation(self):
|
||||
"""Test BM25 term frequency saturation behavior."""
|
||||
# Test with documents where term frequencies can be meaningfully compared
|
||||
documents = [
|
||||
"rare word text", # TF = 1 for "rare"
|
||||
"rare rare word", # TF = 2 for "rare"
|
||||
"rare rare rare rare rare word", # TF = 5 for "rare"
|
||||
"other content", # No "rare" term
|
||||
]
|
||||
query_terms = ["rare"]
|
||||
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
# Documents with the term should have positive scores
|
||||
assert scores[0] > 0.0 # TF=1
|
||||
assert scores[1] > 0.0 # TF=2
|
||||
assert scores[2] > 0.0 # TF=5
|
||||
assert scores[3] == 0.0 # TF=0
|
||||
|
||||
# Scores should increase with term frequency, but with diminishing returns
|
||||
assert scores[1] > scores[0] # TF=2 > TF=1
|
||||
assert scores[2] > scores[1] # TF=5 > TF=2
|
||||
|
||||
# Check that increases demonstrate saturation effect
|
||||
increase_1_to_2 = scores[1] - scores[0]
|
||||
increase_2_to_5 = scores[2] - scores[1]
|
||||
assert increase_1_to_2 > 0
|
||||
assert increase_2_to_5 > 0
|
||||
|
||||
def test_compute_bm25_idf_calculation(self):
|
||||
"""Test BM25 inverse document frequency calculation."""
|
||||
# "rare" appears in only one document, "common" appears in multiple
|
||||
documents = ["rare term", "common term", "common word"]
|
||||
query_terms = ["rare", "common"]
|
||||
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
# First document should have higher score due to rare term having higher IDF
|
||||
assert scores[0] > scores[1] # rare term gets higher IDF
|
||||
assert scores[0] > scores[2]
|
||||
|
||||
def test_compute_bm25_zero_parameters(self, sample_documents, query_terms):
|
||||
"""Test BM25 with edge case parameters."""
|
||||
# Test with k1=0 (no term frequency scaling)
|
||||
scores_k1_zero = compute_bm25(sample_documents, query_terms, k1=0.0)
|
||||
assert len(scores_k1_zero) == len(sample_documents)
|
||||
|
||||
# Test with b=0 (no document length normalization)
|
||||
scores_b_zero = compute_bm25(sample_documents, query_terms, b=0.0)
|
||||
assert len(scores_b_zero) == len(sample_documents)
|
||||
|
||||
# Test with b=1 (full document length normalization)
|
||||
scores_b_one = compute_bm25(sample_documents, query_terms, b=1.0)
|
||||
assert len(scores_b_one) == len(sample_documents)
|
||||
|
||||
def test_tfidf_vs_bm25_comparison(self, sample_documents, query_terms):
|
||||
"""Test that TF-IDF and BM25 produce different but related scores."""
|
||||
tfidf_scores = compute_tfidf(sample_documents, query_terms)
|
||||
bm25_scores = compute_bm25(sample_documents, query_terms)
|
||||
|
||||
# Both should return same number of scores
|
||||
assert len(tfidf_scores) == len(bm25_scores) == len(sample_documents)
|
||||
|
||||
# For documents that match, both should be positive
|
||||
for i in range(len(sample_documents)):
|
||||
if tfidf_scores[i] > 0:
|
||||
assert bm25_scores[i] > 0, f"Document {i} has TF-IDF score but zero BM25 score"
|
||||
if bm25_scores[i] > 0:
|
||||
assert tfidf_scores[i] > 0, f"Document {i} has BM25 score but zero TF-IDF score"
|
||||
|
||||
def test_compute_tfidf_special_characters(self):
|
||||
"""Test TF-IDF with documents containing special characters."""
|
||||
documents = ["hello, world!", "world... hello?", "no match here"]
|
||||
query_terms = ["hello", "world"]
|
||||
|
||||
scores = compute_tfidf(documents, query_terms)
|
||||
|
||||
# Should handle punctuation and still match terms
|
||||
assert len(scores) == 3
|
||||
# Note: Current implementation does simple split(), so punctuation stays attached
|
||||
# This tests the current behavior - may need updating if tokenization improves
|
||||
|
||||
def test_compute_bm25_special_characters(self):
|
||||
"""Test BM25 with documents containing special characters."""
|
||||
documents = ["hello, world!", "world... hello?", "no match here"]
|
||||
query_terms = ["hello", "world"]
|
||||
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
# Should handle punctuation and still match terms
|
||||
assert len(scores) == 3
|
||||
# Same tokenization behavior as TF-IDF
|
||||
|
||||
def test_compute_tfidf_whitespace_handling(self):
|
||||
"""Test TF-IDF with various whitespace scenarios."""
|
||||
documents = [
|
||||
" hello world ", # Extra spaces
|
||||
"\thello\tworld\t", # Tabs
|
||||
"hello\nworld", # Newlines
|
||||
"", # Empty string
|
||||
]
|
||||
query_terms = ["hello", "world"]
|
||||
|
||||
scores = compute_tfidf(documents, query_terms)
|
||||
|
||||
assert len(scores) == 4
|
||||
# First three should have positive scores (they contain the terms)
|
||||
assert scores[0] > 0.0
|
||||
assert scores[1] > 0.0
|
||||
assert scores[2] > 0.0
|
||||
# Last should be zero (empty document)
|
||||
assert scores[3] == 0.0
|
||||
|
||||
def test_compute_bm25_whitespace_handling(self):
|
||||
"""Test BM25 with various whitespace scenarios."""
|
||||
documents = [
|
||||
" hello world ", # Extra spaces
|
||||
"\thello\tworld\t", # Tabs
|
||||
"hello\nworld", # Newlines
|
||||
"", # Empty string
|
||||
]
|
||||
query_terms = ["hello", "world"]
|
||||
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
assert len(scores) == 4
|
||||
# First three should have positive scores (they contain the terms)
|
||||
assert scores[0] > 0.0
|
||||
assert scores[1] > 0.0
|
||||
assert scores[2] > 0.0
|
||||
# Last should be zero (empty document)
|
||||
assert scores[3] == 0.0
|
||||
|
||||
def test_compute_tfidf_mathematical_properties(self):
|
||||
"""Test mathematical properties of TF-IDF scores."""
|
||||
documents = ["cat dog", "cat", "dog"]
|
||||
query_terms = ["cat"]
|
||||
|
||||
scores = compute_tfidf(documents, query_terms)
|
||||
|
||||
# All scores should be non-negative
|
||||
assert all(score >= 0.0 for score in scores)
|
||||
|
||||
# Documents containing the term should have positive scores
|
||||
assert scores[0] > 0.0 # contains "cat"
|
||||
assert scores[1] > 0.0 # contains "cat"
|
||||
assert scores[2] == 0.0 # doesn't contain "cat"
|
||||
|
||||
def test_compute_bm25_mathematical_properties(self):
|
||||
"""Test mathematical properties of BM25 scores."""
|
||||
documents = ["cat dog", "cat", "dog"]
|
||||
query_terms = ["cat"]
|
||||
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
# All scores should be non-negative
|
||||
assert all(score >= 0.0 for score in scores)
|
||||
|
||||
# Documents containing the term should have positive scores
|
||||
assert scores[0] > 0.0 # contains "cat"
|
||||
assert scores[1] > 0.0 # contains "cat"
|
||||
assert scores[2] == 0.0 # doesn't contain "cat"
|
||||
|
||||
def test_compute_tfidf_duplicate_terms_in_query(self):
|
||||
"""Test TF-IDF with duplicate terms in query."""
|
||||
documents = ["cat dog bird", "cat cat dog", "bird bird bird"]
|
||||
query_terms = ["cat", "cat", "dog"] # "cat" appears twice
|
||||
|
||||
scores = compute_tfidf(documents, query_terms)
|
||||
|
||||
# Should handle duplicate query terms gracefully
|
||||
assert len(scores) == 3
|
||||
assert all(isinstance(score, float) for score in scores)
|
||||
|
||||
# First two documents should have positive scores
|
||||
assert scores[0] > 0.0
|
||||
assert scores[1] > 0.0
|
||||
# Third document only contains "bird", so should have zero score
|
||||
assert scores[2] == 0.0
|
||||
|
||||
def test_compute_bm25_duplicate_terms_in_query(self):
|
||||
"""Test BM25 with duplicate terms in query."""
|
||||
documents = ["cat dog bird", "cat cat dog", "bird bird bird"]
|
||||
query_terms = ["cat", "cat", "dog"] # "cat" appears twice
|
||||
|
||||
scores = compute_bm25(documents, query_terms)
|
||||
|
||||
# Should handle duplicate query terms gracefully
|
||||
assert len(scores) == 3
|
||||
assert all(isinstance(score, float) for score in scores)
|
||||
|
||||
# First two documents should have positive scores
|
||||
assert scores[0] > 0.0
|
||||
assert scores[1] > 0.0
|
||||
# Third document only contains "bird", so should have zero score
|
||||
assert scores[2] == 0.0
|
||||
392
src/backend/tests/unit/components/data/test_kb_ingest.py
Normal file
392
src/backend/tests/unit/components/data/test_kb_ingest.py
Normal file
|
|
@ -0,0 +1,392 @@
|
|||
import json
|
||||
from pathlib import Path
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pandas as pd
|
||||
import pytest
|
||||
from langflow.components.data.kb_ingest import KBIngestionComponent
|
||||
from langflow.schema.data import Data
|
||||
|
||||
from tests.base import ComponentTestBaseWithoutClient
|
||||
|
||||
|
||||
class TestKBIngestionComponent(ComponentTestBaseWithoutClient):
|
||||
@pytest.fixture
|
||||
def component_class(self):
|
||||
"""Return the component class to test."""
|
||||
return KBIngestionComponent
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def mock_knowledge_base_path(self, tmp_path):
|
||||
"""Mock the knowledge base root path directly."""
|
||||
with patch("langflow.components.data.kb_ingest.KNOWLEDGE_BASES_ROOT_PATH", tmp_path):
|
||||
yield
|
||||
|
||||
@pytest.fixture
|
||||
def default_kwargs(self, tmp_path):
|
||||
"""Return default kwargs for component instantiation."""
|
||||
# Create a sample DataFrame
|
||||
data_df = pd.DataFrame(
|
||||
{"text": ["Sample text 1", "Sample text 2"], "title": ["Title 1", "Title 2"], "category": ["cat1", "cat2"]}
|
||||
)
|
||||
|
||||
# Create column configuration
|
||||
column_config = [
|
||||
{"column_name": "text", "vectorize": True, "identifier": False},
|
||||
{"column_name": "title", "vectorize": False, "identifier": False},
|
||||
{"column_name": "category", "vectorize": False, "identifier": True},
|
||||
]
|
||||
|
||||
# Create knowledge base directory
|
||||
kb_name = "test_kb"
|
||||
kb_path = tmp_path / kb_name
|
||||
kb_path.mkdir(exist_ok=True)
|
||||
|
||||
# Create embedding metadata file
|
||||
metadata = {
|
||||
"embedding_provider": "HuggingFace",
|
||||
"embedding_model": "sentence-transformers/all-MiniLM-L6-v2",
|
||||
"api_key": None,
|
||||
"api_key_used": False,
|
||||
"chunk_size": 1000,
|
||||
"created_at": "2024-01-01T00:00:00Z",
|
||||
}
|
||||
(kb_path / "embedding_metadata.json").write_text(json.dumps(metadata))
|
||||
|
||||
return {
|
||||
"knowledge_base": kb_name,
|
||||
"input_df": data_df,
|
||||
"column_config": column_config,
|
||||
"chunk_size": 1000,
|
||||
"kb_root_path": str(tmp_path),
|
||||
"api_key": None,
|
||||
"allow_duplicates": False,
|
||||
"silent_errors": False,
|
||||
}
|
||||
|
||||
@pytest.fixture
|
||||
def file_names_mapping(self):
|
||||
"""Return file names mapping for version testing."""
|
||||
# This is a new component, so it doesn't exist in older versions
|
||||
return []
|
||||
|
||||
def test_validate_column_config_valid(self, component_class, default_kwargs):
|
||||
"""Test column configuration validation with valid config."""
|
||||
component = component_class(**default_kwargs)
|
||||
data_df = default_kwargs["input_df"]
|
||||
|
||||
config_list = component._validate_column_config(data_df)
|
||||
|
||||
assert len(config_list) == 3
|
||||
assert config_list[0]["column_name"] == "text"
|
||||
assert config_list[0]["vectorize"] is True
|
||||
|
||||
def test_validate_column_config_invalid_column(self, component_class, default_kwargs):
|
||||
"""Test column configuration validation with invalid column name."""
|
||||
# Modify column config to include non-existent column
|
||||
invalid_config = [{"column_name": "nonexistent", "vectorize": True, "identifier": False}]
|
||||
default_kwargs["column_config"] = invalid_config
|
||||
|
||||
component = component_class(**default_kwargs)
|
||||
data_df = default_kwargs["input_df"]
|
||||
|
||||
with pytest.raises(ValueError, match="Column 'nonexistent' not found in DataFrame"):
|
||||
component._validate_column_config(data_df)
|
||||
|
||||
def test_validate_column_config_silent_errors(self, component_class, default_kwargs):
|
||||
"""Test column configuration validation with silent errors enabled."""
|
||||
# Modify column config to include non-existent column
|
||||
invalid_config = [{"column_name": "nonexistent", "vectorize": True, "identifier": False}]
|
||||
default_kwargs["column_config"] = invalid_config
|
||||
default_kwargs["silent_errors"] = True
|
||||
|
||||
component = component_class(**default_kwargs)
|
||||
data_df = default_kwargs["input_df"]
|
||||
|
||||
# Should not raise exception with silent_errors=True
|
||||
config_list = component._validate_column_config(data_df)
|
||||
assert isinstance(config_list, list)
|
||||
|
||||
def test_get_embedding_provider(self, component_class, default_kwargs):
|
||||
"""Test embedding provider detection."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Test OpenAI provider
|
||||
assert component._get_embedding_provider("text-embedding-ada-002") == "OpenAI"
|
||||
|
||||
# Test HuggingFace provider
|
||||
assert component._get_embedding_provider("sentence-transformers/all-MiniLM-L6-v2") == "HuggingFace"
|
||||
|
||||
# Test Cohere provider
|
||||
assert component._get_embedding_provider("embed-english-v3.0") == "Cohere"
|
||||
|
||||
# Test custom provider
|
||||
assert component._get_embedding_provider("custom-model") == "Custom"
|
||||
|
||||
@patch("langchain_huggingface.HuggingFaceEmbeddings")
|
||||
def test_build_embeddings_huggingface(self, mock_hf_embeddings, component_class, default_kwargs):
|
||||
"""Test building HuggingFace embeddings."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
mock_embeddings = MagicMock()
|
||||
mock_hf_embeddings.return_value = mock_embeddings
|
||||
|
||||
result = component._build_embeddings("sentence-transformers/all-MiniLM-L6-v2", None)
|
||||
|
||||
mock_hf_embeddings.assert_called_once_with(model="sentence-transformers/all-MiniLM-L6-v2")
|
||||
assert result == mock_embeddings
|
||||
|
||||
@patch("langchain_openai.OpenAIEmbeddings")
|
||||
def test_build_embeddings_openai(self, mock_openai_embeddings, component_class, default_kwargs):
|
||||
"""Test building OpenAI embeddings."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
mock_embeddings = MagicMock()
|
||||
mock_openai_embeddings.return_value = mock_embeddings
|
||||
|
||||
result = component._build_embeddings("text-embedding-ada-002", "test-api-key")
|
||||
|
||||
mock_openai_embeddings.assert_called_once_with(
|
||||
model="text-embedding-ada-002", api_key="test-api-key", chunk_size=1000
|
||||
)
|
||||
assert result == mock_embeddings
|
||||
|
||||
def test_build_embeddings_openai_no_key(self, component_class, default_kwargs):
|
||||
"""Test building OpenAI embeddings without API key raises error."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
with pytest.raises(ValueError, match="OpenAI API key is required"):
|
||||
component._build_embeddings("text-embedding-ada-002", None)
|
||||
|
||||
@patch("langchain_cohere.CohereEmbeddings")
|
||||
def test_build_embeddings_cohere(self, mock_cohere_embeddings, component_class, default_kwargs):
|
||||
"""Test building Cohere embeddings."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
mock_embeddings = MagicMock()
|
||||
mock_cohere_embeddings.return_value = mock_embeddings
|
||||
|
||||
result = component._build_embeddings("embed-english-v3.0", "test-api-key")
|
||||
|
||||
mock_cohere_embeddings.assert_called_once_with(model="embed-english-v3.0", cohere_api_key="test-api-key")
|
||||
assert result == mock_embeddings
|
||||
|
||||
def test_build_embeddings_cohere_no_key(self, component_class, default_kwargs):
|
||||
"""Test building Cohere embeddings without API key raises error."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
with pytest.raises(ValueError, match="Cohere API key is required"):
|
||||
component._build_embeddings("embed-english-v3.0", None)
|
||||
|
||||
def test_build_embeddings_custom_not_supported(self, component_class, default_kwargs):
|
||||
"""Test building custom embeddings raises NotImplementedError."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
with pytest.raises(NotImplementedError, match="Custom embedding models not yet supported"):
|
||||
component._build_embeddings("custom-model", "test-key")
|
||||
|
||||
@patch("langflow.components.data.kb_ingest.get_settings_service")
|
||||
@patch("langflow.components.data.kb_ingest.encrypt_api_key")
|
||||
def test_build_embedding_metadata(self, mock_encrypt, mock_get_settings, component_class, default_kwargs):
|
||||
"""Test building embedding metadata."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
mock_settings = MagicMock()
|
||||
mock_get_settings.return_value = mock_settings
|
||||
mock_encrypt.return_value = "encrypted_key"
|
||||
|
||||
metadata = component._build_embedding_metadata("sentence-transformers/all-MiniLM-L6-v2", "test-key")
|
||||
|
||||
assert metadata["embedding_provider"] == "HuggingFace"
|
||||
assert metadata["embedding_model"] == "sentence-transformers/all-MiniLM-L6-v2"
|
||||
assert metadata["api_key"] == "encrypted_key"
|
||||
assert metadata["api_key_used"] is True
|
||||
assert metadata["chunk_size"] == 1000
|
||||
assert "created_at" in metadata
|
||||
|
||||
def test_build_column_metadata(self, component_class, default_kwargs):
|
||||
"""Test building column metadata."""
|
||||
component = component_class(**default_kwargs)
|
||||
data_df = default_kwargs["input_df"]
|
||||
config_list = default_kwargs["column_config"]
|
||||
|
||||
metadata = component._build_column_metadata(config_list, data_df)
|
||||
|
||||
assert metadata["total_columns"] == 3
|
||||
assert metadata["mapped_columns"] == 3
|
||||
assert metadata["unmapped_columns"] == 0
|
||||
assert len(metadata["columns"]) == 3
|
||||
assert "text" in metadata["summary"]["vectorized_columns"]
|
||||
assert "category" in metadata["summary"]["identifier_columns"]
|
||||
|
||||
def test_convert_df_to_data_objects(self, component_class, default_kwargs):
|
||||
"""Test converting DataFrame to Data objects."""
|
||||
component = component_class(**default_kwargs)
|
||||
data_df = default_kwargs["input_df"]
|
||||
config_list = default_kwargs["column_config"]
|
||||
|
||||
# Mock Chroma to avoid actual vector store operations
|
||||
with patch("langflow.components.data.kb_ingest.Chroma") as mock_chroma:
|
||||
mock_chroma_instance = MagicMock()
|
||||
mock_chroma_instance.get.return_value = {"metadatas": []}
|
||||
mock_chroma.return_value = mock_chroma_instance
|
||||
|
||||
data_objects = component._convert_df_to_data_objects(data_df, config_list)
|
||||
|
||||
assert len(data_objects) == 2
|
||||
assert all(isinstance(obj, Data) for obj in data_objects)
|
||||
|
||||
# Check first data object
|
||||
first_obj = data_objects[0]
|
||||
assert "text" in first_obj.data
|
||||
assert "title" in first_obj.data
|
||||
assert "category" in first_obj.data
|
||||
assert "_id" in first_obj.data
|
||||
|
||||
def test_convert_df_to_data_objects_no_duplicates(self, component_class, default_kwargs):
|
||||
"""Test converting DataFrame to Data objects with duplicate prevention."""
|
||||
default_kwargs["allow_duplicates"] = False
|
||||
component = component_class(**default_kwargs)
|
||||
data_df = default_kwargs["input_df"]
|
||||
config_list = default_kwargs["column_config"]
|
||||
|
||||
# Mock Chroma with existing hash
|
||||
with patch("langflow.components.data.kb_ingest.Chroma") as mock_chroma:
|
||||
# Simulate existing document with same hash
|
||||
existing_hash = "some_existing_hash"
|
||||
mock_chroma_instance = MagicMock()
|
||||
mock_chroma_instance.get.return_value = {"metadatas": [{"_id": existing_hash}]}
|
||||
mock_chroma.return_value = mock_chroma_instance
|
||||
|
||||
# Mock hashlib to return the existing hash for first row
|
||||
with patch("langflow.components.data.kb_ingest.hashlib.sha256") as mock_hash:
|
||||
mock_hash_obj = MagicMock()
|
||||
mock_hash_obj.hexdigest.side_effect = [existing_hash, "different_hash"]
|
||||
mock_hash.return_value = mock_hash_obj
|
||||
|
||||
data_objects = component._convert_df_to_data_objects(data_df, config_list)
|
||||
|
||||
# Should only return one object (second row) since first is duplicate
|
||||
assert len(data_objects) == 1
|
||||
|
||||
def test_is_valid_collection_name(self, component_class, default_kwargs):
|
||||
"""Test collection name validation."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Valid names
|
||||
assert component.is_valid_collection_name("valid_name") is True
|
||||
assert component.is_valid_collection_name("valid-name") is True
|
||||
assert component.is_valid_collection_name("ValidName123") is True
|
||||
|
||||
# Invalid names
|
||||
assert component.is_valid_collection_name("ab") is False # Too short
|
||||
assert component.is_valid_collection_name("a" * 64) is False # Too long
|
||||
assert component.is_valid_collection_name("_invalid") is False # Starts with underscore
|
||||
assert component.is_valid_collection_name("invalid_") is False # Ends with underscore
|
||||
assert component.is_valid_collection_name("invalid@name") is False # Invalid character
|
||||
|
||||
@patch("langflow.components.data.kb_ingest.json.loads")
|
||||
@patch("langflow.components.data.kb_ingest.decrypt_api_key")
|
||||
def test_build_kb_info_success(self, mock_decrypt, mock_json_loads, component_class, default_kwargs):
|
||||
"""Test successful KB info building."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Mock metadata loading
|
||||
mock_json_loads.return_value = {
|
||||
"embedding_model": "sentence-transformers/all-MiniLM-L6-v2",
|
||||
"api_key": "encrypted_key",
|
||||
}
|
||||
mock_decrypt.return_value = "decrypted_key"
|
||||
|
||||
# Mock vector store creation
|
||||
with patch.object(component, "_create_vector_store"), patch.object(component, "_save_kb_files"):
|
||||
result = component.build_kb_info()
|
||||
|
||||
assert isinstance(result, Data)
|
||||
assert "kb_id" in result.data
|
||||
assert "kb_name" in result.data
|
||||
assert "rows" in result.data
|
||||
assert result.data["rows"] == 2
|
||||
|
||||
def test_build_kb_info_with_silent_errors(self, component_class, default_kwargs):
|
||||
"""Test KB info building with silent errors enabled."""
|
||||
default_kwargs["silent_errors"] = True
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Remove the metadata file to cause an error
|
||||
kb_path = Path(default_kwargs["kb_root_path"]) / default_kwargs["knowledge_base"]
|
||||
metadata_file = kb_path / "embedding_metadata.json"
|
||||
if metadata_file.exists():
|
||||
metadata_file.unlink()
|
||||
|
||||
# Should not raise exception with silent_errors=True
|
||||
result = component.build_kb_info()
|
||||
assert isinstance(result, Data)
|
||||
assert "error" in result.data
|
||||
|
||||
def test_get_knowledge_bases(self, component_class, default_kwargs, tmp_path):
|
||||
"""Test getting list of knowledge bases."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Create additional test directories
|
||||
(tmp_path / "kb1").mkdir()
|
||||
(tmp_path / "kb2").mkdir()
|
||||
(tmp_path / ".hidden").mkdir() # Should be ignored
|
||||
|
||||
kb_list = component._get_knowledge_bases()
|
||||
|
||||
assert "test_kb" in kb_list
|
||||
assert "kb1" in kb_list
|
||||
assert "kb2" in kb_list
|
||||
assert ".hidden" not in kb_list
|
||||
|
||||
@patch("langflow.components.data.kb_ingest.Path.exists")
|
||||
def test_get_knowledge_bases_no_path(self, mock_exists, component_class, default_kwargs):
|
||||
"""Test getting knowledge bases when path doesn't exist."""
|
||||
component = component_class(**default_kwargs)
|
||||
mock_exists.return_value = False
|
||||
|
||||
kb_list = component._get_knowledge_bases()
|
||||
assert kb_list == []
|
||||
|
||||
def test_update_build_config_new_kb(self, component_class, default_kwargs):
|
||||
"""Test updating build config for new knowledge base creation."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
build_config = {"knowledge_base": {"value": None, "options": []}}
|
||||
|
||||
field_value = {
|
||||
"01_new_kb_name": "new_test_kb",
|
||||
"02_embedding_model": "sentence-transformers/all-MiniLM-L6-v2",
|
||||
"03_api_key": None,
|
||||
}
|
||||
|
||||
# Mock embedding validation
|
||||
with (
|
||||
patch.object(component, "_build_embeddings") as mock_build_emb,
|
||||
patch.object(component, "_save_embedding_metadata"),
|
||||
patch.object(component, "_get_knowledge_bases") as mock_get_kbs,
|
||||
):
|
||||
mock_embeddings = MagicMock()
|
||||
mock_embeddings.embed_query.return_value = [0.1, 0.2, 0.3]
|
||||
mock_build_emb.return_value = mock_embeddings
|
||||
mock_get_kbs.return_value = ["new_test_kb"]
|
||||
|
||||
result = component.update_build_config(build_config, field_value, "knowledge_base")
|
||||
|
||||
assert result["knowledge_base"]["value"] == "new_test_kb"
|
||||
assert "new_test_kb" in result["knowledge_base"]["options"]
|
||||
|
||||
def test_update_build_config_invalid_kb_name(self, component_class, default_kwargs):
|
||||
"""Test updating build config with invalid KB name."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
build_config = {"knowledge_base": {"value": None, "options": []}}
|
||||
field_value = {
|
||||
"01_new_kb_name": "invalid@name", # Invalid character
|
||||
"02_embedding_model": "sentence-transformers/all-MiniLM-L6-v2",
|
||||
"03_api_key": None,
|
||||
}
|
||||
|
||||
with pytest.raises(ValueError, match="Invalid knowledge base name"):
|
||||
component.update_build_config(build_config, field_value, "knowledge_base")
|
||||
368
src/backend/tests/unit/components/data/test_kb_retrieval.py
Normal file
368
src/backend/tests/unit/components/data/test_kb_retrieval.py
Normal file
|
|
@ -0,0 +1,368 @@
|
|||
import contextlib
|
||||
import json
|
||||
from pathlib import Path
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
from langflow.components.data.kb_retrieval import KBRetrievalComponent
|
||||
|
||||
from tests.base import ComponentTestBaseWithoutClient
|
||||
|
||||
|
||||
class TestKBRetrievalComponent(ComponentTestBaseWithoutClient):
|
||||
@pytest.fixture
|
||||
def component_class(self):
|
||||
"""Return the component class to test."""
|
||||
return KBRetrievalComponent
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def mock_knowledge_base_path(self, tmp_path):
|
||||
"""Mock the knowledge base root path directly."""
|
||||
with patch("langflow.components.data.kb_retrieval.KNOWLEDGE_BASES_ROOT_PATH", tmp_path):
|
||||
yield
|
||||
|
||||
@pytest.fixture
|
||||
def default_kwargs(self, tmp_path):
|
||||
"""Return default kwargs for component instantiation."""
|
||||
# Create knowledge base directory structure
|
||||
kb_name = "test_kb"
|
||||
kb_path = tmp_path / kb_name
|
||||
kb_path.mkdir(exist_ok=True)
|
||||
|
||||
# Create embedding metadata file
|
||||
metadata = {
|
||||
"embedding_provider": "HuggingFace",
|
||||
"embedding_model": "sentence-transformers/all-MiniLM-L6-v2",
|
||||
"api_key": None,
|
||||
"api_key_used": False,
|
||||
"chunk_size": 1000,
|
||||
"created_at": "2024-01-01T00:00:00Z",
|
||||
}
|
||||
(kb_path / "embedding_metadata.json").write_text(json.dumps(metadata))
|
||||
|
||||
return {
|
||||
"knowledge_base": kb_name,
|
||||
"kb_root_path": str(tmp_path),
|
||||
"api_key": None,
|
||||
"search_query": "",
|
||||
"top_k": 5,
|
||||
"include_embeddings": True,
|
||||
}
|
||||
|
||||
@pytest.fixture
|
||||
def file_names_mapping(self):
|
||||
"""Return file names mapping for version testing."""
|
||||
# This is a new component, so it doesn't exist in older versions
|
||||
return []
|
||||
|
||||
def test_get_knowledge_bases(self, component_class, default_kwargs, tmp_path):
|
||||
"""Test getting list of knowledge bases."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Create additional test directories
|
||||
(tmp_path / "kb1").mkdir()
|
||||
(tmp_path / "kb2").mkdir()
|
||||
(tmp_path / ".hidden").mkdir() # Should be ignored
|
||||
|
||||
kb_list = component._get_knowledge_bases()
|
||||
|
||||
assert "test_kb" in kb_list
|
||||
assert "kb1" in kb_list
|
||||
assert "kb2" in kb_list
|
||||
assert ".hidden" not in kb_list
|
||||
|
||||
@patch("langflow.components.data.kb_retrieval.Path.exists")
|
||||
def test_get_knowledge_bases_no_path(self, mock_exists, component_class, default_kwargs):
|
||||
"""Test getting knowledge bases when path doesn't exist."""
|
||||
component = component_class(**default_kwargs)
|
||||
mock_exists.return_value = False
|
||||
|
||||
kb_list = component._get_knowledge_bases()
|
||||
assert kb_list == []
|
||||
|
||||
def test_update_build_config(self, component_class, default_kwargs, tmp_path):
|
||||
"""Test updating build configuration."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Create additional KB directories
|
||||
(tmp_path / "kb1").mkdir()
|
||||
(tmp_path / "kb2").mkdir()
|
||||
|
||||
build_config = {"knowledge_base": {"value": "test_kb", "options": []}}
|
||||
|
||||
result = component.update_build_config(build_config, None, "knowledge_base")
|
||||
|
||||
assert "test_kb" in result["knowledge_base"]["options"]
|
||||
assert "kb1" in result["knowledge_base"]["options"]
|
||||
assert "kb2" in result["knowledge_base"]["options"]
|
||||
|
||||
def test_update_build_config_invalid_kb(self, component_class, default_kwargs):
|
||||
"""Test updating build config when selected KB is not available."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
build_config = {"knowledge_base": {"value": "nonexistent_kb", "options": ["test_kb"]}}
|
||||
|
||||
result = component.update_build_config(build_config, None, "knowledge_base")
|
||||
|
||||
assert result["knowledge_base"]["value"] is None
|
||||
|
||||
def test_get_kb_metadata_success(self, component_class, default_kwargs):
|
||||
"""Test successful metadata loading."""
|
||||
component = component_class(**default_kwargs)
|
||||
kb_path = Path(default_kwargs["kb_root_path"]) / default_kwargs["knowledge_base"]
|
||||
|
||||
with patch("langflow.components.data.kb_retrieval.decrypt_api_key") as mock_decrypt:
|
||||
mock_decrypt.return_value = "decrypted_key"
|
||||
|
||||
metadata = component._get_kb_metadata(kb_path)
|
||||
|
||||
assert metadata["embedding_provider"] == "HuggingFace"
|
||||
assert metadata["embedding_model"] == "sentence-transformers/all-MiniLM-L6-v2"
|
||||
assert "chunk_size" in metadata
|
||||
|
||||
def test_get_kb_metadata_no_file(self, component_class, default_kwargs, tmp_path):
|
||||
"""Test metadata loading when file doesn't exist."""
|
||||
component = component_class(**default_kwargs)
|
||||
nonexistent_path = tmp_path / "nonexistent"
|
||||
nonexistent_path.mkdir()
|
||||
|
||||
metadata = component._get_kb_metadata(nonexistent_path)
|
||||
|
||||
assert metadata == {}
|
||||
|
||||
def test_get_kb_metadata_json_error(self, component_class, default_kwargs, tmp_path):
|
||||
"""Test metadata loading with invalid JSON."""
|
||||
component = component_class(**default_kwargs)
|
||||
kb_path = tmp_path / "invalid_json_kb"
|
||||
kb_path.mkdir()
|
||||
|
||||
# Create invalid JSON file
|
||||
(kb_path / "embedding_metadata.json").write_text("invalid json content")
|
||||
|
||||
metadata = component._get_kb_metadata(kb_path)
|
||||
|
||||
assert metadata == {}
|
||||
|
||||
def test_get_kb_metadata_decrypt_error(self, component_class, default_kwargs, tmp_path):
|
||||
"""Test metadata loading with decryption error."""
|
||||
component = component_class(**default_kwargs)
|
||||
kb_path = tmp_path / "decrypt_error_kb"
|
||||
kb_path.mkdir()
|
||||
|
||||
# Create metadata with encrypted key
|
||||
metadata = {
|
||||
"embedding_provider": "OpenAI",
|
||||
"embedding_model": "text-embedding-ada-002",
|
||||
"api_key": "encrypted_key",
|
||||
"chunk_size": 1000,
|
||||
}
|
||||
(kb_path / "embedding_metadata.json").write_text(json.dumps(metadata))
|
||||
|
||||
with patch("langflow.components.data.kb_retrieval.decrypt_api_key") as mock_decrypt:
|
||||
mock_decrypt.side_effect = ValueError("Decryption failed")
|
||||
|
||||
result = component._get_kb_metadata(kb_path)
|
||||
|
||||
assert result["api_key"] is None
|
||||
|
||||
@patch("langchain_huggingface.HuggingFaceEmbeddings")
|
||||
def test_build_embeddings_huggingface(self, mock_hf_embeddings, component_class, default_kwargs):
|
||||
"""Test building HuggingFace embeddings."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
metadata = {
|
||||
"embedding_provider": "HuggingFace",
|
||||
"embedding_model": "sentence-transformers/all-MiniLM-L6-v2",
|
||||
"chunk_size": 1000,
|
||||
}
|
||||
|
||||
mock_embeddings = MagicMock()
|
||||
mock_hf_embeddings.return_value = mock_embeddings
|
||||
|
||||
result = component._build_embeddings(metadata)
|
||||
|
||||
mock_hf_embeddings.assert_called_once_with(model="sentence-transformers/all-MiniLM-L6-v2")
|
||||
assert result == mock_embeddings
|
||||
|
||||
@patch("langchain_openai.OpenAIEmbeddings")
|
||||
def test_build_embeddings_openai(self, mock_openai_embeddings, component_class, default_kwargs):
|
||||
"""Test building OpenAI embeddings."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
metadata = {
|
||||
"embedding_provider": "OpenAI",
|
||||
"embedding_model": "text-embedding-ada-002",
|
||||
"api_key": "test-api-key",
|
||||
"chunk_size": 1000,
|
||||
}
|
||||
|
||||
mock_embeddings = MagicMock()
|
||||
mock_openai_embeddings.return_value = mock_embeddings
|
||||
|
||||
result = component._build_embeddings(metadata)
|
||||
|
||||
mock_openai_embeddings.assert_called_once_with(
|
||||
model="text-embedding-ada-002", api_key="test-api-key", chunk_size=1000
|
||||
)
|
||||
assert result == mock_embeddings
|
||||
|
||||
def test_build_embeddings_openai_no_key(self, component_class, default_kwargs):
|
||||
"""Test building OpenAI embeddings without API key raises error."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
metadata = {
|
||||
"embedding_provider": "OpenAI",
|
||||
"embedding_model": "text-embedding-ada-002",
|
||||
"api_key": None,
|
||||
"chunk_size": 1000,
|
||||
}
|
||||
|
||||
with pytest.raises(ValueError, match="OpenAI API key is required"):
|
||||
component._build_embeddings(metadata)
|
||||
|
||||
@patch("langchain_cohere.CohereEmbeddings")
|
||||
def test_build_embeddings_cohere(self, mock_cohere_embeddings, component_class, default_kwargs):
|
||||
"""Test building Cohere embeddings."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
metadata = {
|
||||
"embedding_provider": "Cohere",
|
||||
"embedding_model": "embed-english-v3.0",
|
||||
"api_key": "test-api-key",
|
||||
"chunk_size": 1000,
|
||||
}
|
||||
|
||||
mock_embeddings = MagicMock()
|
||||
mock_cohere_embeddings.return_value = mock_embeddings
|
||||
|
||||
result = component._build_embeddings(metadata)
|
||||
|
||||
mock_cohere_embeddings.assert_called_once_with(model="embed-english-v3.0", cohere_api_key="test-api-key")
|
||||
assert result == mock_embeddings
|
||||
|
||||
def test_build_embeddings_cohere_no_key(self, component_class, default_kwargs):
|
||||
"""Test building Cohere embeddings without API key raises error."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
metadata = {
|
||||
"embedding_provider": "Cohere",
|
||||
"embedding_model": "embed-english-v3.0",
|
||||
"api_key": None,
|
||||
"chunk_size": 1000,
|
||||
}
|
||||
|
||||
with pytest.raises(ValueError, match="Cohere API key is required"):
|
||||
component._build_embeddings(metadata)
|
||||
|
||||
def test_build_embeddings_custom_not_supported(self, component_class, default_kwargs):
|
||||
"""Test building custom embeddings raises NotImplementedError."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
metadata = {"embedding_provider": "Custom", "embedding_model": "custom-model", "api_key": "test-key"}
|
||||
|
||||
with pytest.raises(NotImplementedError, match="Custom embedding models not yet supported"):
|
||||
component._build_embeddings(metadata)
|
||||
|
||||
def test_build_embeddings_unsupported_provider(self, component_class, default_kwargs):
|
||||
"""Test building embeddings with unsupported provider raises NotImplementedError."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
metadata = {"embedding_provider": "UnsupportedProvider", "embedding_model": "some-model", "api_key": "test-key"}
|
||||
|
||||
with pytest.raises(NotImplementedError, match="Embedding provider 'UnsupportedProvider' is not supported"):
|
||||
component._build_embeddings(metadata)
|
||||
|
||||
def test_build_embeddings_with_user_api_key(self, component_class, default_kwargs):
|
||||
"""Test that user-provided API key overrides stored one."""
|
||||
# Create a mock secret input
|
||||
|
||||
mock_secret = MagicMock()
|
||||
mock_secret.get_secret_value.return_value = "user-provided-key"
|
||||
|
||||
default_kwargs["api_key"] = mock_secret
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
metadata = {
|
||||
"embedding_provider": "OpenAI",
|
||||
"embedding_model": "text-embedding-ada-002",
|
||||
"api_key": "stored-key",
|
||||
"chunk_size": 1000,
|
||||
}
|
||||
|
||||
with patch("langchain_openai.OpenAIEmbeddings") as mock_openai:
|
||||
mock_embeddings = MagicMock()
|
||||
mock_openai.return_value = mock_embeddings
|
||||
|
||||
component._build_embeddings(metadata)
|
||||
|
||||
mock_openai.assert_called_once_with(
|
||||
model="text-embedding-ada-002", api_key="user-provided-key", chunk_size=1000
|
||||
)
|
||||
|
||||
def test_get_chroma_kb_data_no_metadata(self, component_class, default_kwargs, tmp_path):
|
||||
"""Test retrieving data when metadata is missing."""
|
||||
# Remove metadata file
|
||||
kb_path = tmp_path / default_kwargs["knowledge_base"]
|
||||
metadata_file = kb_path / "embedding_metadata.json"
|
||||
if metadata_file.exists():
|
||||
metadata_file.unlink()
|
||||
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
with pytest.raises(ValueError, match="Metadata not found for knowledge base"):
|
||||
component.get_chroma_kb_data()
|
||||
|
||||
def test_get_chroma_kb_data_path_construction(self, component_class, default_kwargs):
|
||||
"""Test that get_chroma_kb_data constructs the correct paths."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Test that the component correctly builds the KB path
|
||||
|
||||
assert component.kb_root_path == default_kwargs["kb_root_path"]
|
||||
assert component.knowledge_base == default_kwargs["knowledge_base"]
|
||||
|
||||
# Test that paths are correctly expanded
|
||||
expanded_path = Path(component.kb_root_path).expanduser()
|
||||
assert expanded_path.exists() # tmp_path should exist
|
||||
|
||||
# Verify method exists with correct parameters
|
||||
assert hasattr(component, "get_chroma_kb_data")
|
||||
assert hasattr(component, "search_query")
|
||||
assert hasattr(component, "top_k")
|
||||
assert hasattr(component, "include_embeddings")
|
||||
|
||||
def test_get_chroma_kb_data_method_exists(self, component_class, default_kwargs):
|
||||
"""Test that get_chroma_kb_data method exists and can be called."""
|
||||
component = component_class(**default_kwargs)
|
||||
|
||||
# Just verify the method exists and has the right signature
|
||||
assert hasattr(component, "get_chroma_kb_data"), "Component should have get_chroma_kb_data method"
|
||||
|
||||
# Mock all external calls to avoid integration issues
|
||||
with (
|
||||
patch.object(component, "_get_kb_metadata") as mock_get_metadata,
|
||||
patch.object(component, "_build_embeddings") as mock_build_embeddings,
|
||||
patch("langchain_chroma.Chroma"),
|
||||
):
|
||||
mock_get_metadata.return_value = {"embedding_provider": "HuggingFace", "embedding_model": "test-model"}
|
||||
mock_build_embeddings.return_value = MagicMock()
|
||||
|
||||
# This is a unit test focused on the component's internal logic
|
||||
with contextlib.suppress(Exception):
|
||||
component.get_chroma_kb_data()
|
||||
|
||||
# Verify internal methods were called
|
||||
mock_get_metadata.assert_called_once()
|
||||
mock_build_embeddings.assert_called_once()
|
||||
|
||||
def test_include_embeddings_parameter(self, component_class, default_kwargs):
|
||||
"""Test that include_embeddings parameter is properly set."""
|
||||
# Test with embeddings enabled
|
||||
default_kwargs["include_embeddings"] = True
|
||||
component = component_class(**default_kwargs)
|
||||
assert component.include_embeddings is True
|
||||
|
||||
# Test with embeddings disabled
|
||||
default_kwargs["include_embeddings"] = False
|
||||
component = component_class(**default_kwargs)
|
||||
assert component.include_embeddings is False
|
||||
Loading…
Add table
Add a link
Reference in a new issue