From a1769bc93b95cda1f54c437b5a313d2547376f2a Mon Sep 17 00:00:00 2001 From: Ishaan Jaffer Date: Mon, 26 Jan 2026 15:58:29 -0800 Subject: [PATCH] _save_vector_store_to_db_from_rag_ingest --- litellm/proxy/rag_endpoints/endpoints.py | 105 ++++++++++++++++++++++- 1 file changed, 102 insertions(+), 3 deletions(-) diff --git a/litellm/proxy/rag_endpoints/endpoints.py b/litellm/proxy/rag_endpoints/endpoints.py index f9e0b724cb4..21d598c5495 100644 --- a/litellm/proxy/rag_endpoints/endpoints.py +++ b/litellm/proxy/rag_endpoints/endpoints.py @@ -26,11 +26,65 @@ from litellm.proxy.common_utils.http_parsing_utils import ( router = APIRouter() +def _build_file_metadata_entry( + response: Any, + file_data: Optional[Tuple[str, bytes, str]] = None, + file_url: Optional[str] = None, +) -> Dict[str, Any]: + """ + Build a file metadata entry for storing in vector_store_metadata. + + Args: + response: The response from litellm.aingest containing file_id + file_data: Optional tuple of (filename, content, content_type) + file_url: Optional URL if file was ingested from URL + + Returns: + Dictionary with file metadata (file_id, filename, file_url, ingested_at, etc.) + """ + from datetime import datetime, timezone + + # Extract file_id from response + file_id = None + if hasattr(response, "get"): + file_id = response.get("file_id") + elif hasattr(response, "file_id"): + file_id = response.file_id + + # Extract file information from file_data tuple + filename = None + file_size = None + content_type = None + + if file_data: + filename = file_data[0] + file_size = len(file_data[1]) if len(file_data) > 1 else None + content_type = file_data[2] if len(file_data) > 2 else None + + # Build file metadata entry + file_entry = { + "file_id": file_id, + "filename": filename, + "file_url": file_url, + "ingested_at": datetime.now(timezone.utc).isoformat(), + } + + # Add optional fields if available + if file_size is not None: + file_entry["file_size"] = file_size + if content_type is not None: + file_entry["content_type"] = content_type + + return file_entry + + async def _save_vector_store_to_db_from_rag_ingest( response: Any, ingest_options: Dict[str, Any], prisma_client, user_api_key_dict: UserAPIKeyAuth, + file_data: Optional[Tuple[str, bytes, str]] = None, + file_url: Optional[str] = None, ) -> None: """ Helper function to save a newly created vector store from RAG ingest to the database. @@ -70,6 +124,18 @@ async def _save_vector_store_to_db_from_rag_ingest( vector_store_config = ingest_options.get("vector_store", {}) custom_llm_provider = vector_store_config.get("custom_llm_provider") + + # Extract litellm_vector_store_params for custom name and description + litellm_vector_store_params = ingest_options.get("litellm_vector_store_params", {}) + custom_vector_store_name = litellm_vector_store_params.get("vector_store_name") + custom_vector_store_description = litellm_vector_store_params.get("vector_store_description") + + # Build file metadata entry using helper + file_entry = _build_file_metadata_entry( + response=response, + file_data=file_data, + file_url=file_url, + ) try: # Check if vector store already exists in database @@ -85,12 +151,22 @@ async def _save_vector_store_to_db_from_rag_ingest( f"Saving newly created vector store {vector_store_id} to database" ) + # Initialize metadata with first file + initial_metadata = { + "ingested_files": [file_entry] + } + + # Use custom name if provided, otherwise default + vector_store_name = custom_vector_store_name or f"RAG Vector Store - {vector_store_id[:8]}" + vector_store_description = custom_vector_store_description or "Created via RAG ingest endpoint" + await create_vector_store_in_db( vector_store_id=vector_store_id, custom_llm_provider=custom_llm_provider or "openai", prisma_client=prisma_client, - vector_store_name=f"RAG Vector Store - {vector_store_id[:8]}", - vector_store_description="Created via RAG ingest endpoint", + vector_store_name=vector_store_name, + vector_store_description=vector_store_description, + vector_store_metadata=initial_metadata, ) verbose_proxy_logger.info( @@ -98,7 +174,28 @@ async def _save_vector_store_to_db_from_rag_ingest( ) else: verbose_proxy_logger.info( - f"Vector store {vector_store_id} already exists in database, skipping creation" + f"Vector store {vector_store_id} already exists, appending file to metadata" + ) + + # Update existing vector store with new file + existing_metadata = existing_vector_store.vector_store_metadata or {} + if isinstance(existing_metadata, str): + import json + existing_metadata = json.loads(existing_metadata) + + ingested_files = existing_metadata.get("ingested_files", []) + ingested_files.append(file_entry) + existing_metadata["ingested_files"] = ingested_files + + # Update the vector store + from litellm.proxy.utils import safe_dumps + await prisma_client.db.litellm_managedvectorstorestable.update( + where={"vector_store_id": vector_store_id}, + data={"vector_store_metadata": safe_dumps(existing_metadata)} + ) + + verbose_proxy_logger.info( + f"Added file {file_entry.get('filename') or file_entry.get('file_url', 'Unknown')} to vector store {vector_store_id} metadata" ) except Exception as db_error: # Log the error but don't fail the request since ingestion succeeded @@ -283,6 +380,8 @@ async def rag_ingest( ingest_options=ingest_options, prisma_client=prisma_client, user_api_key_dict=user_api_key_dict, + file_data=file_data, + file_url=file_url, ) else: verbose_proxy_logger.warning(