Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
76 changes: 76 additions & 0 deletions backend/open_webui/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -1382,6 +1382,82 @@ def reachable(host: str, port: int) -> bool:
)


# --- BEGIN EXTERNAL RETRIEVAL PATCH ---
# External retrieval engine: allows delegating document search to an external HTTP service
RAG_RETRIEVAL_ENGINE = ConfigVar(
"RAG_RETRIEVAL_ENGINE",
"rag.retrieval_engine",
os.environ.get("RAG_RETRIEVAL_ENGINE", ""),
)

RAG_EXTERNAL_RETRIEVAL_URL = ConfigVar(
"RAG_EXTERNAL_RETRIEVAL_URL",
"rag.external_retrieval_url",
os.environ.get("RAG_EXTERNAL_RETRIEVAL_URL", ""),
)

RAG_EXTERNAL_RETRIEVAL_API_KEY = ConfigVar(
"RAG_EXTERNAL_RETRIEVAL_API_KEY",
"rag.external_retrieval_api_key",
os.environ.get("RAG_EXTERNAL_RETRIEVAL_API_KEY", ""),
)

RAG_EXTERNAL_RETRIEVAL_TIMEOUT = ConfigVar(
"RAG_EXTERNAL_RETRIEVAL_TIMEOUT",
"rag.external_retrieval_timeout",
os.environ.get("RAG_EXTERNAL_RETRIEVAL_TIMEOUT", ""),
)

RAG_EXTERNAL_BYPASS_QUERY_GENERATION = ConfigVar(
"RAG_EXTERNAL_BYPASS_QUERY_GENERATION",
"rag.external_bypass_query_generation",
os.environ.get("RAG_EXTERNAL_BYPASS_QUERY_GENERATION", "false").lower() == "true",
)

RAG_EXTERNAL_MESSAGE_COUNT = ConfigVar(
"RAG_EXTERNAL_MESSAGE_COUNT",
"rag.external_message_count",
int(os.environ.get("RAG_EXTERNAL_MESSAGE_COUNT", "10")),
)

RAG_EXTERNAL_USER_MESSAGES_ONLY = ConfigVar(
"RAG_EXTERNAL_USER_MESSAGES_ONLY",
"rag.external_user_messages_only",
os.environ.get("RAG_EXTERNAL_USER_MESSAGES_ONLY", "false").lower() == "true",
)
# --- END EXTERNAL RETRIEVAL PATCH ---


# --- BEGIN EXTERNAL INGESTION PATCH ---
# External ingestion engine: delegates document chunking, embedding, and vector
# storage to an external HTTP service instead of running save_docs_to_vector_db
# in-process. Default off; set EXTERNAL_INGESTION_ENGINE=external to enable.
EXTERNAL_INGESTION_ENGINE = ConfigVar(
"EXTERNAL_INGESTION_ENGINE",
"rag.external_ingestion_engine",
os.environ.get("EXTERNAL_INGESTION_ENGINE", ""),
)

EXTERNAL_INGESTION_URL = ConfigVar(
"EXTERNAL_INGESTION_URL",
"rag.external_ingestion_url",
os.environ.get("EXTERNAL_INGESTION_URL", ""),
)

EXTERNAL_INGESTION_API_KEY = ConfigVar(
"EXTERNAL_INGESTION_API_KEY",
"rag.external_ingestion_api_key",
os.environ.get("EXTERNAL_INGESTION_API_KEY", ""),
)

EXTERNAL_INGESTION_TIMEOUT = ConfigVar(
"EXTERNAL_INGESTION_TIMEOUT",
"rag.external_ingestion_timeout",
os.environ.get("EXTERNAL_INGESTION_TIMEOUT", "300"),
)
# --- END EXTERNAL INGESTION PATCH ---


RAG_TEXT_SPLITTER = ConfigVar(
'RAG_TEXT_SPLITTER',
'rag.text_splitter',
Expand Down
30 changes: 30 additions & 0 deletions backend/open_webui/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -341,6 +341,21 @@
RAG_OPENAI_API_BASE_URL,
RAG_OPENAI_API_KEY,
RAG_RELEVANCE_THRESHOLD,
# --- BEGIN EXTERNAL RETRIEVAL PATCH ---
RAG_RETRIEVAL_ENGINE,
RAG_EXTERNAL_RETRIEVAL_URL,
RAG_EXTERNAL_RETRIEVAL_API_KEY,
RAG_EXTERNAL_RETRIEVAL_TIMEOUT,
RAG_EXTERNAL_BYPASS_QUERY_GENERATION,
RAG_EXTERNAL_MESSAGE_COUNT,
RAG_EXTERNAL_USER_MESSAGES_ONLY,
# --- END EXTERNAL RETRIEVAL PATCH ---
# --- BEGIN EXTERNAL INGESTION PATCH ---
EXTERNAL_INGESTION_ENGINE,
EXTERNAL_INGESTION_URL,
EXTERNAL_INGESTION_API_KEY,
EXTERNAL_INGESTION_TIMEOUT,
# --- END EXTERNAL INGESTION PATCH ---
RAG_RERANKING_BATCH_SIZE,
RAG_RERANKING_ENGINE,
RAG_RERANKING_MODEL,
Expand Down Expand Up @@ -1057,6 +1072,21 @@ async def lifespan(app: FastAPI):
app.state.config.RAG_EXTERNAL_RERANKER_API_KEY = RAG_EXTERNAL_RERANKER_API_KEY
app.state.config.RAG_EXTERNAL_RERANKER_TIMEOUT = RAG_EXTERNAL_RERANKER_TIMEOUT
app.state.config.RAG_RERANKING_BATCH_SIZE = RAG_RERANKING_BATCH_SIZE
# --- BEGIN EXTERNAL RETRIEVAL PATCH ---
app.state.config.RAG_RETRIEVAL_ENGINE = RAG_RETRIEVAL_ENGINE
app.state.config.RAG_EXTERNAL_RETRIEVAL_URL = RAG_EXTERNAL_RETRIEVAL_URL
app.state.config.RAG_EXTERNAL_RETRIEVAL_API_KEY = RAG_EXTERNAL_RETRIEVAL_API_KEY
app.state.config.RAG_EXTERNAL_RETRIEVAL_TIMEOUT = RAG_EXTERNAL_RETRIEVAL_TIMEOUT
app.state.config.RAG_EXTERNAL_BYPASS_QUERY_GENERATION = RAG_EXTERNAL_BYPASS_QUERY_GENERATION
app.state.config.RAG_EXTERNAL_MESSAGE_COUNT = RAG_EXTERNAL_MESSAGE_COUNT
app.state.config.RAG_EXTERNAL_USER_MESSAGES_ONLY = RAG_EXTERNAL_USER_MESSAGES_ONLY
# --- END EXTERNAL RETRIEVAL PATCH ---
# --- BEGIN EXTERNAL INGESTION PATCH ---
app.state.config.EXTERNAL_INGESTION_ENGINE = EXTERNAL_INGESTION_ENGINE
app.state.config.EXTERNAL_INGESTION_URL = EXTERNAL_INGESTION_URL
app.state.config.EXTERNAL_INGESTION_API_KEY = EXTERNAL_INGESTION_API_KEY
app.state.config.EXTERNAL_INGESTION_TIMEOUT = EXTERNAL_INGESTION_TIMEOUT
# --- END EXTERNAL INGESTION PATCH ---

app.state.config.RAG_TEMPLATE = RAG_TEMPLATE

Expand Down
210 changes: 210 additions & 0 deletions backend/open_webui/retrieval/external.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,210 @@
# --- BEGIN EXTERNAL RETRIEVAL PATCH ---
# External retrieval engine: delegates document search to an external HTTP service
# instead of querying the built-in vector DB directly.
# --- END EXTERNAL RETRIEVAL PATCH ---
# --- BEGIN EXTERNAL INGESTION PATCH ---
# External ingestion engine: delegates document chunking, embedding, and vector
# storage to an external HTTP service instead of running save_docs_to_vector_db
# in-process.
# --- END EXTERNAL INGESTION PATCH ---

import logging
from typing import Optional, List

import requests

from open_webui.env import ENABLE_FORWARD_USER_INFO_HEADERS, REQUESTS_VERIFY
from open_webui.utils.headers import include_user_info_headers

log = logging.getLogger(__name__)


def query_external_retrieval(
url: str,
api_key: str,
queries: List[str],
collection_names: List[str],
k: int,
timeout: Optional[str] = None,
user=None,
messages: Optional[List[dict]] = None,
) -> Optional[dict]:
"""
Query an external retrieval service.

POST {url}/search with queries + collection_names + k.
Optionally includes the chat messages so the external service can
extract/generate its own queries from the conversation. Open WebUI's
QUERY_GENERATION_PROMPT_TEMPLATE is intentionally NOT forwarded — the
external service runs its own query generation.
Returns dict with keys: documents, metadatas, distances (matching internal format).
Returns None on error.
"""
payload = {
"queries": queries,
"collection_names": collection_names,
"k": k,
}

if messages is not None:
payload["messages"] = messages

try:
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {api_key}",
}

if ENABLE_FORWARD_USER_INFO_HEADERS and user:
headers = include_user_info_headers(headers, user)

request_timeout = int(timeout) if timeout else None

log.info(
f"query_external_retrieval: url={url}, queries={queries}, "
f"messages={len(messages) if messages else 0}, "
f"collections={collection_names}, k={k}"
)

r = requests.post(
f"{url}/search",
headers=headers,
json=payload,
timeout=request_timeout,
verify=REQUESTS_VERIFY,
)
r.raise_for_status()
data = r.json()

if "documents" in data:
return data
else:
log.error("No documents found in external retrieval response")
return None

except Exception as e:
log.exception(f"Error in external retrieval: {e}")
return None


# --- BEGIN EXTERNAL INGESTION PATCH ---
def process_file_external_ingestion(
url: str,
api_key: str,
file_id: str,
filename: str,
collection_name: str,
user_id: str,
local_file_path: Optional[str] = None,
s3_bucket: Optional[str] = None,
s3_key: Optional[str] = None,
timeout: Optional[int] = 300,
) -> Optional[dict]:
"""
Delegate document ingestion to an external HTTP service.

PUT {url}/api/v1/ingest with either:
- JSON body containing s3_bucket + s3_key (preferred when storage is S3)
- multipart body containing the file (fallback when storage is local)

Returns the service's response dict on success
({"status": True, "collection_name": ..., "chunks_count": ...})
or None on transport / server error. Caller treats None as failure.
"""
try:
headers = {"Authorization": f"Bearer {api_key}"}
endpoint = f"{url.rstrip('/')}/api/v1/ingest"

if s3_bucket and s3_key:
payload = {
"s3_bucket": s3_bucket,
"s3_key": s3_key,
"file_id": file_id,
"filename": filename,
"collection_name": collection_name,
"collection_type": "file",
"user_id": user_id,
"overwrite": True,
}
log.info(
f"process_file_external_ingestion (s3): file_id={file_id}, "
f"collection={collection_name}, s3={s3_bucket}/{s3_key}"
)
r = requests.put(
endpoint,
headers={**headers, "Content-Type": "application/json"},
json=payload,
timeout=timeout,
verify=REQUESTS_VERIFY,
)
elif local_file_path:
log.info(
f"process_file_external_ingestion (multipart): file_id={file_id}, "
f"collection={collection_name}, path={local_file_path}"
)
with open(local_file_path, "rb") as fh:
files = {"file": (filename, fh)}
data = {
"file_id": file_id,
"filename": filename,
"collection_name": collection_name,
"collection_type": "file",
"user_id": user_id,
"overwrite": "true",
}
r = requests.put(
endpoint,
headers=headers,
files=files,
data=data,
timeout=timeout,
verify=REQUESTS_VERIFY,
)
else:
log.error(
"process_file_external_ingestion: no S3 reference and no local file path"
)
return None

r.raise_for_status()
return r.json()

except Exception as e:
log.exception(f"Error in external ingestion: {e}")
return None


def delete_file_external_ingestion(
url: str,
api_key: str,
file_id: str,
timeout: int = 300,
) -> Optional[dict]:
"""Tell the external ingestion service to drop a file's vectors.

DELETE {url}/api/v1/documents/{file_id}. The vector store lives behind the
external service now, so Open WebUI's own vector-DB cleanup no longer reaches
it — this call keeps the two in sync when a file is deleted.

Best-effort: returns the service's response dict on success, or None on any
transport/server error (logged, never raised). Callers MUST NOT let a failed
cleanup block the user's file deletion. ``file_id`` is the bare file UUID —
the service matches on meta.file_id, not the "file-" collection name.
"""
try:
headers = {"Authorization": f"Bearer {api_key}"}
endpoint = f"{url.rstrip('/')}/api/v1/documents/{file_id}"
log.info(f"delete_file_external_ingestion: file_id={file_id}")
r = requests.delete(
endpoint,
headers=headers,
timeout=timeout,
verify=REQUESTS_VERIFY,
)
r.raise_for_status()
return r.json()

except Exception as e:
log.exception(f"Error in external ingestion delete: {e}")
return None
# --- END EXTERNAL INGESTION PATCH ---
25 changes: 25 additions & 0 deletions backend/open_webui/retrieval/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,11 @@
)
from langchain_community.retrievers import BM25Retriever
from langchain_core.documents import Document

# --- BEGIN EXTERNAL RETRIEVAL PATCH ---
from open_webui.retrieval.external import query_external_retrieval
# --- END EXTERNAL RETRIEVAL PATCH ---

from open_webui.config import (
RAG_EMBEDDING_CONTENT_PREFIX,
RAG_EMBEDDING_PREFIX_FIELD_NAME,
Expand Down Expand Up @@ -1162,6 +1167,9 @@ async def get_sources_from_items(
hybrid_search,
full_context=False,
user: UserModel | None = None,
# --- BEGIN EXTERNAL RETRIEVAL PATCH ---
messages: Optional[list] = None,
# --- END EXTERNAL RETRIEVAL PATCH ---
):
log.debug(f'items: {items} {queries} {embedding_function} {reranking_function} {full_context}')

Expand Down Expand Up @@ -1410,6 +1418,23 @@ async def get_sources_from_items(
# Sync helper makes blocking VECTOR_DB_CLIENT calls;
# offload so the async caller's event loop stays free.
query_result = await asyncio.to_thread(get_all_items_from_collections, collection_names)
# --- BEGIN EXTERNAL RETRIEVAL PATCH ---
elif request.app.state.config.RAG_RETRIEVAL_ENGINE == "external":
# The external retrieval service runs its own query
# generation, so Open WebUI's QUERY_GENERATION_PROMPT_TEMPLATE
# is intentionally NOT forwarded.
query_result = await asyncio.to_thread(
query_external_retrieval,
url=request.app.state.config.RAG_EXTERNAL_RETRIEVAL_URL,
api_key=request.app.state.config.RAG_EXTERNAL_RETRIEVAL_API_KEY,
queries=queries,
collection_names=list(collection_names),
k=k,
timeout=request.app.state.config.RAG_EXTERNAL_RETRIEVAL_TIMEOUT,
user=user,
messages=messages,
)
# --- END EXTERNAL RETRIEVAL PATCH ---
else:
query_result = await query_collection(
request,
Expand Down
Loading