diff --git a/surfsense_backend/alembic/versions/22_add_elasticsearch_connector_enum.py b/surfsense_backend/alembic/versions/23_add_elasticsearch_connector_enum.py
similarity index 88%
rename from surfsense_backend/alembic/versions/22_add_elasticsearch_connector_enum.py
rename to surfsense_backend/alembic/versions/23_add_elasticsearch_connector_enum.py
index ad13af1f8..cb2725b00 100644
--- a/surfsense_backend/alembic/versions/22_add_elasticsearch_connector_enum.py
+++ b/surfsense_backend/alembic/versions/23_add_elasticsearch_connector_enum.py
@@ -1,7 +1,7 @@
"""Add ElasticSearch connector enums
-Revision ID: 22
-Revises: 21
+Revision ID: 23
+Revises: 22
Create Date: 2025-10-08 12:00:00.000000
"""
@@ -11,8 +11,8 @@ from collections.abc import Sequence
from alembic import op
# revision identifiers
-revision: str = "22"
-down_revision: str | None = "21"
+revision: str = "23"
+down_revision: str | None = "22"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
@@ -52,4 +52,5 @@ def upgrade() -> None:
def downgrade() -> None:
+ """Remove 'ELASTICSEARCH_CONNECTOR' from enum types."""
pass
diff --git a/surfsense_backend/app/connectors/elasticsearch_connector.py b/surfsense_backend/app/connectors/elasticsearch_connector.py
index 4a5ceb2b1..19921472c 100644
--- a/surfsense_backend/app/connectors/elasticsearch_connector.py
+++ b/surfsense_backend/app/connectors/elasticsearch_connector.py
@@ -1,256 +1,95 @@
-import json
-import logging
-from collections.abc import AsyncGenerator
+from __future__ import annotations
+
+from collections.abc import Iterable
from typing import Any
from elasticsearch import AsyncElasticsearch
-try:
- from elasticsearch.exceptions import AuthenticationException
-except ImportError:
- AuthenticationException = Exception
-
-try:
- from elastic_transport import (
- ApiError,
- ConnectionError as TransportConnectionError,
- TlsError,
- )
-except ImportError:
- # Fallback for older versions
- TransportConnectionError = ConnectionError
- TlsError = Exception
- ApiError = Exception
-
-
-from app.db import DocumentType
-
-logger = logging.getLogger(__name__)
-
-
-# Create a simple document structure for Elasticsearch results
-class ElasticsearchDocument:
- def __init__(
- self, title: str, content: str, url: str, metadata: dict | None = None
- ):
- self.title = title
- self.content = content
- self.url = url
- self.metadata = metadata or {}
- self.document_type = DocumentType.ELASTICSEARCH_CONNECTOR
-
class ElasticsearchConnector:
- def __init__(self, config: dict):
- self.config = config
- self.client: AsyncElasticsearch | None = None
+ """Async helper for executing Elasticsearch searches."""
- async def connect(self) -> bool:
- """Establish connection to Elasticsearch"""
- try:
- # Build the URL with proper scheme
- scheme = "https" if self.config.get("ssl_enabled", True) else "http"
- port = self.config.get("port", 443 if scheme == "https" else 9200)
- hostname = self.config["hostname"]
+ def __init__(self, config: dict[str, Any]) -> None:
+ self._config = config
+ self._client: AsyncElasticsearch | None = None
- # Construct the full URL
- url = f"{scheme}://{hostname}:{port}"
+ async def __aenter__(self) -> ElasticsearchConnector:
+ self._client = AsyncElasticsearch(**self._build_client_kwargs())
+ return self
- connection_params = {
- "hosts": [url],
- "verify_certs": self.config.get("ssl_enabled", True),
- "request_timeout": 30,
- }
+ async def __aexit__(self, exc_type, exc, tb) -> None:
+ if self._client is not None:
+ await self._client.close()
+ self._client = None
- # Handle different authentication methods
- auth_method = self.config.get("auth_method", "none")
+ def _build_client_kwargs(self) -> dict[str, Any]:
+ hostname = self._config.get("hostname")
+ if not hostname:
+ raise ValueError("Elasticsearch connector config missing 'hostname'.")
- if (
- auth_method == "basic"
- and self.config.get("username")
- and self.config.get("password")
- ):
- connection_params["basic_auth"] = (
- self.config["username"],
- self.config["password"],
+ port = int(self._config.get("port", 9200))
+ ssl_enabled = bool(self._config.get("ssl_enabled", True))
+ scheme = "https" if ssl_enabled else "http"
+
+ kwargs: dict[str, Any] = {
+ "hosts": [{"host": hostname, "port": port, "scheme": scheme}],
+ "verify_certs": ssl_enabled,
+ }
+
+ auth_method = self._config.get("auth_method", "api_key")
+ if auth_method == "api_key":
+ api_key = self._config.get("ELASTICSEARCH_API_KEY")
+ if not api_key:
+ raise ValueError(
+ "Elasticsearch connector config missing 'ELASTICSEARCH_API_KEY'."
)
- elif auth_method == "api_key" and self.config.get("api_key"):
- api_key = self.config["api_key"]
- if api_key.startswith("ApiKey "):
- connection_params["headers"] = {"Authorization": api_key}
- else:
- connection_params["api_key"] = api_key
+ kwargs["api_key"] = api_key
+ elif auth_method == "basic":
+ username = self._config.get("username")
+ password = self._config.get("password")
+ if not username or not password:
+ raise ValueError(
+ "Elasticsearch basic auth requires 'username' and 'password'."
+ )
+ kwargs["basic_auth"] = (username, password)
+ else:
+ raise ValueError(f"Unsupported Elasticsearch auth method '{auth_method}'.")
- logger.info(f"Connecting to Elasticsearch at {url}")
- self.client = AsyncElasticsearch(**connection_params)
+ if not ssl_enabled:
+ kwargs["verify_certs"] = False
+ kwargs["ssl_show_warn"] = False
- # Test connection
- info = await self.client.info()
- logger.info(
- f"Successfully connected to Elasticsearch cluster: {info.get('cluster_name', 'Unknown')}"
- )
- return True
-
- except (ConnectionError, TransportConnectionError, TlsError, ApiError) as e:
- logger.error(f"Failed to connect to Elasticsearch: {e}")
- return False
- except Exception as e:
- logger.error(f"Unexpected error connecting to Elasticsearch: {e}")
- return False
-
- async def disconnect(self):
- """Close Elasticsearch connection"""
- if self.client:
- await self.client.close()
- self.client = None
-
- async def test_connection(self) -> dict[str, Any]:
- """Test the Elasticsearch connection and return cluster info"""
- try:
- if not self.client:
- await self.connect()
-
- if not self.client:
- return {"success": False, "error": "Failed to establish connection"}
-
- info = await self.client.info()
- indices = await self.client.cat.indices(format="json")
-
- return {
- "success": True,
- "cluster_name": info.get("cluster_name"),
- "version": info.get("version", {}).get("number"),
- "indices_count": len(indices) if indices else 0,
- "indices": [
- idx["index"] for idx in indices if not idx["index"].startswith(".")
- ],
- }
- except Exception as e:
- return {"success": False, "error": str(e)}
-
- async def get_indices(self) -> list[str]:
- """Get list of available indices"""
- try:
- if not self.client:
- await self.connect()
-
- indices = await self.client.cat.indices(format="json")
- return [idx["index"] for idx in indices if not idx["index"].startswith(".")]
- except Exception as e:
- logger.error(f"Error fetching indices: {e}")
- return []
+ return kwargs
async def search_documents(
self,
+ *,
+ indices: Iterable[str] | None,
query: str,
- indices: list[str] | None = None,
- size: int = 100,
- fields: list[str] | None = None,
- ) -> AsyncGenerator[ElasticsearchDocument, None]:
- """Search documents in Elasticsearch and yield them as ElasticsearchDocument objects"""
- try:
- if not self.client:
- await self.connect()
-
- if not self.client:
- logger.error("No Elasticsearch client available")
- return
-
- # Build search query
- search_query = {
- "multi_match": {
- "query": query,
- "fields": fields or ["*"],
- "type": "best_fields",
- "fuzziness": "AUTO",
- }
- }
-
- # Search across specified indices or all indices
- index_pattern = ",".join(indices) if indices else "*"
-
- response = await self.client.search(
- index=index_pattern,
- query=search_query,
- size=size,
- source=True,
+ search_fields: Iterable[str] | None,
+ max_documents: int,
+ ) -> list[dict[str, Any]]:
+ if self._client is None:
+ raise RuntimeError(
+ "Elasticsearch client not initialised. Use as an async context manager."
)
- for hit in response["hits"]["hits"]:
- try:
- # Convert Elasticsearch document to ElasticsearchDocument
- document = await self._convert_es_document(hit)
- if document:
- yield document
- except Exception as e:
- logger.warning(f"Failed to convert document {hit.get('_id')}: {e}")
- continue
+ size = max(1, min(int(max_documents or 1000), 10_000))
+ index_param = ",".join(indices) if indices else None
- except Exception as e:
- logger.error(f"Error searching Elasticsearch: {e}")
+ if not query or query == "*":
+ query_body: dict[str, Any] = {"match_all": {}}
+ else:
+ query_body = {"query_string": {"query": query}}
+ if search_fields:
+ fields = [field for field in search_fields if field]
+ if fields:
+ query_body["query_string"]["fields"] = fields
- async def get_document_by_id(
- self, index: str, doc_id: str
- ) -> ElasticsearchDocument | None:
- """Get a specific document by ID"""
- try:
- if not self.client:
- await self.connect()
-
- response = await self.client.get(index=index, id=doc_id)
- return await self._convert_es_document(response)
- except Exception as e:
- logger.error(f"Error fetching document {doc_id} from {index}: {e}")
- return None
-
- async def _convert_es_document(
- self, es_doc: dict[str, Any]
- ) -> ElasticsearchDocument | None:
- """Convert Elasticsearch document to ElasticsearchDocument"""
- try:
- source = es_doc.get("_source", {})
-
- # Extract title - try common field names
- title = (
- source.get("title")
- or source.get("name")
- or source.get("subject")
- or source.get("filename")
- or f"Document {es_doc.get('_id', 'Unknown')}"
- )
-
- # Extract content - try common field names
- content = (
- source.get("content")
- or source.get("text")
- or source.get("body")
- or source.get("message")
- or json.dumps(source) # Fallback to full source as JSON
- )
-
- # Generate URL pointing back to Elasticsearch
- url = f"elasticsearch://{es_doc.get('_index')}/{es_doc.get('_id')}"
-
- # Create metadata
- metadata = {
- "elasticsearch_index": es_doc.get("_index"),
- "elasticsearch_id": es_doc.get("_id"),
- "elasticsearch_score": es_doc.get("_score"),
- "source_fields": list(source.keys()),
- **{
- k: v
- for k, v in source.items()
- if k not in ["content", "text", "body", "message"]
- and isinstance(v, str | int | float | bool)
- },
- }
-
- return ElasticsearchDocument(
- title=title[:500], # Limit title length
- content=content,
- url=url,
- metadata=metadata,
- )
- except Exception as e:
- logger.error(f"Error converting Elasticsearch document: {e}")
- return None
+ response = await self._client.search(
+ index=index_param,
+ size=size,
+ query=query_body,
+ _source_includes=list(search_fields) if search_fields else None,
+ )
+ return response.get("hits", {}).get("hits", [])
diff --git a/surfsense_backend/app/schemas/elasticsearch.py b/surfsense_backend/app/schemas/elasticsearch.py
index d392d4a20..c5cb29d87 100644
--- a/surfsense_backend/app/schemas/elasticsearch.py
+++ b/surfsense_backend/app/schemas/elasticsearch.py
@@ -12,7 +12,9 @@ class ElasticsearchConnectorConfig(BaseModel):
password: str | None = Field(
default=None, description="Password for authentication"
)
- api_key: str | None = Field(default=None, description="API key for authentication")
+ ELASTICSEARCH_API_KEY: str | None = Field(
+ default=None, description="API key for authentication"
+ )
ssl_enabled: bool = Field(default=True, description="Whether to use SSL/TLS")
indices: list[str] | None = Field(
default=None, description="Specific indices to search (optional)"
@@ -33,7 +35,7 @@ class ElasticsearchTestConnectionRequest(BaseModel):
port: int = 9200
username: str | None = None
password: str | None = None
- api_key: str | None = None
+ ELASTICSEARCH_API_KEY: str | None = None
ssl_enabled: bool = True
diff --git a/surfsense_backend/app/schemas/search_source_connector.py b/surfsense_backend/app/schemas/search_source_connector.py
index f5dbb0e3d..32f3e1ea0 100644
--- a/surfsense_backend/app/schemas/search_source_connector.py
+++ b/surfsense_backend/app/schemas/search_source_connector.py
@@ -215,7 +215,7 @@ class SearchSourceConnectorBase(BaseModel):
"port",
"username",
"password",
- "api_key",
+ "ELASTICSEARCH_API_KEY",
"auth_method",
"ssl_enabled",
"indices",
@@ -237,7 +237,7 @@ class SearchSourceConnectorBase(BaseModel):
if auth_method == "basic":
if not config.get("username") or not config.get("password"):
raise ValueError("Username and password required for basic auth")
- elif auth_method == "api_key" and not config.get("api_key"):
+ elif auth_method == "api_key" and not config.get("ELASTICSEARCH_API_KEY"):
raise ValueError("API key required for api_key auth method")
# Validate that all config keys are allowed
diff --git a/surfsense_backend/app/tasks/connector_indexers/elasticsearch_indexer.py b/surfsense_backend/app/tasks/connector_indexers/elasticsearch_indexer.py
index 5fbe409ad..1f2d6dcb0 100644
--- a/surfsense_backend/app/tasks/connector_indexers/elasticsearch_indexer.py
+++ b/surfsense_backend/app/tasks/connector_indexers/elasticsearch_indexer.py
@@ -1,76 +1,43 @@
-"""
-Elasticsearch connector indexer.
-"""
+from __future__ import annotations
-import asyncio
-from datetime import datetime
+import json
+from collections.abc import Iterable, Sequence
+from dataclasses import dataclass
+from datetime import UTC, datetime
+from hashlib import sha256
+from typing import Any
-from sqlalchemy.exc import SQLAlchemyError
+from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
-from app.config import config
from app.connectors.elasticsearch_connector import ElasticsearchConnector
-from app.db import Document, DocumentType, SearchSourceConnectorType
-from app.services.llm_service import get_user_long_context_llm
-from app.services.task_logging_service import TaskLoggingService
-from app.utils.document_converters import (
- create_document_chunks,
- generate_content_hash,
- generate_document_summary,
-)
-
-from .base import (
- check_duplicate_document_by_hash,
- get_connector_by_id,
- logger,
- update_connector_last_indexed,
-)
+from app.db import Chunk, Document, DocumentType, SearchSourceConnector
-class ElasticsearchIndexer:
- """Indexer class for Elasticsearch data sources (kept for backward compatibility)"""
+@dataclass(slots=True)
+class ElasticsearchSearchParams:
+ indices: Sequence[str] | None
+ query: str
+ search_fields: Sequence[str] | None
+ max_documents: int
+ start_date: str | None = None
+ end_date: str | None = None
- def __init__(self, connector_config: dict):
- self.connector_config = connector_config
- self.connector = ElasticsearchConnector(self.connector_config)
- async def get_documents(self):
- """Get documents from Elasticsearch (kept for backward compatibility)"""
- try:
- if not await self.connector.connect():
- logger.error("Failed to connect to Elasticsearch")
- return
+def _coerce_sequence(value: Any) -> list[str] | None:
+ if value is None:
+ return None
+ if isinstance(value, str):
+ items = [item.strip() for item in value.split(",") if item.strip()]
+ return items or None
+ if isinstance(value, Iterable):
+ items = [str(item).strip() for item in value if str(item).strip()]
+ return items or None
+ return None
- query = self.connector_config.get("query", "*")
- indices = self.connector_config.get("indices", None)
- max_docs = self.connector_config.get("max_documents", 1000)
- fields = self.connector_config.get("search_fields", None)
- logger.info(f"Searching Elasticsearch with query: {query}")
-
- doc_count = 0
- async for document in self.connector.search_documents(
- query=query, indices=indices, size=max_docs, fields=fields
- ):
- if doc_count >= max_docs:
- break
-
- yield document
- doc_count += 1
-
- if doc_count % 10 == 0:
- await asyncio.sleep(0.1)
-
- logger.info(f"Fetched {doc_count} documents from Elasticsearch")
-
- except Exception as e:
- logger.error(f"Error fetching documents from Elasticsearch: {e}")
- finally:
- await self.connector.disconnect()
-
- async def test_connection(self):
- """Test the Elasticsearch connection (kept for backward compatibility)"""
- return await self.connector.test_connection()
+def _format_content(payload: dict[str, Any]) -> str:
+ return json.dumps(payload, ensure_ascii=False, indent=2, sort_keys=True)
async def index_elasticsearch_documents(
@@ -78,310 +45,98 @@ async def index_elasticsearch_documents(
connector_id: int,
search_space_id: int,
user_id: str,
- start_date: str
- | None = None, # Not used for Elasticsearch but kept for consistency
- end_date: str | None = None, # Not used for Elasticsearch but kept for consistency
+ start_date: str | None,
+ end_date: str | None,
+ *,
update_last_indexed: bool = True,
) -> tuple[int, str | None]:
- """
- Index Elasticsearch documents.
+ connector = await session.get(SearchSourceConnector, connector_id)
+ if connector is None:
+ return 0, f"Connector {connector_id} not found."
- Args:
- session: Database session
- connector_id: ID of the Elasticsearch connector
- search_space_id: ID of the search space to store documents in
- user_id: User ID
- start_date: Start date for indexing (not used for Elasticsearch but kept for consistency)
- end_date: End date for indexing (not used for Elasticsearch but kept for consistency)
- update_last_indexed: Whether to update the last_indexed_at timestamp (default: True)
-
- Returns:
- Tuple containing (number of documents indexed, error message or None)
- """
- task_logger = TaskLoggingService(session, search_space_id)
-
- # Log task start
- log_entry = await task_logger.log_task_start(
- task_name="elasticsearch_documents_indexing",
- source="connector_indexing_task",
- message=f"Starting Elasticsearch documents indexing for connector {connector_id}",
- metadata={
- "connector_id": connector_id,
- "user_id": str(user_id),
- "start_date": start_date,
- "end_date": end_date,
- },
+ config = connector.config or {}
+ params = ElasticsearchSearchParams(
+ indices=_coerce_sequence(config.get("indices")),
+ query=str(config.get("query") or "*"),
+ search_fields=_coerce_sequence(config.get("search_fields")),
+ max_documents=int(config.get("max_documents") or 1000),
+ start_date=start_date,
+ end_date=end_date,
)
try:
- # Get the connector
- await task_logger.log_task_progress(
- log_entry,
- f"Retrieving Elasticsearch connector {connector_id} from database",
- {"stage": "connector_retrieval"},
- )
-
- # Get the connector from the database
- connector = await get_connector_by_id(
- session, connector_id, SearchSourceConnectorType.ELASTICSEARCH_CONNECTOR
- )
-
- if not connector:
- await task_logger.log_task_failure(
- log_entry,
- f"Connector with ID {connector_id} not found or is not an Elasticsearch connector",
- "Connector not found",
- {"error_type": "ConnectorNotFound"},
- )
- return (
- 0,
- f"Connector with ID {connector_id} not found or is not an Elasticsearch connector",
+ async with ElasticsearchConnector(config) as elastic:
+ hits = await elastic.search_documents(
+ indices=params.indices,
+ query=params.query,
+ search_fields=params.search_fields,
+ max_documents=params.max_documents,
)
+ except Exception as exc:
+ return 0, f"Elasticsearch query failed: {exc!s}"
- logger.info(f"Starting Elasticsearch indexing for connector {connector_id}")
+ documents_processed = 0
+ ingested_at = datetime.now(UTC).isoformat()
- # Initialize Elasticsearch connector
- await task_logger.log_task_progress(
- log_entry,
- f"Initializing Elasticsearch connector for connector {connector_id}",
- {"stage": "connector_initialization"},
- )
+ try:
+ for hit in hits:
+ source = hit.get("_source") or {}
+ content = _format_content(source)
+ content_hash = sha256(
+ f"{hit.get('_index', '')}::{hit.get('_id', '')}::{content}".encode()
+ ).hexdigest()
- try:
- elasticsearch_connector = ElasticsearchConnector(connector.config)
- except Exception as e:
- await task_logger.log_task_failure(
- log_entry,
- f"Invalid Elasticsearch configuration for connector {connector_id}: {e}",
- "Invalid configuration",
- {"error_type": "InvalidConfiguration"},
- )
- return 0, f"Invalid Elasticsearch configuration: {e}"
-
- # Test connection first
- await task_logger.log_task_progress(
- log_entry,
- f"Testing Elasticsearch connection for connector {connector_id}",
- {"stage": "connection_test"},
- )
-
- connection_test = await elasticsearch_connector.test_connection()
- if not connection_test.get("success"):
- await task_logger.log_task_failure(
- log_entry,
- f"Connection test failed: {connection_test.get('error')}",
- "Connection failed",
- {"error_type": "ConnectionError"},
- )
- return 0, f"Connection test failed: {connection_test.get('error')}"
-
- # Get configuration parameters
- query = connector.config.get("query", "*")
- indices = connector.config.get("indices", None)
- max_docs = connector.config.get("max_documents", 1000)
- fields = connector.config.get("search_fields", None)
-
- await task_logger.log_task_progress(
- log_entry,
- f"Fetching Elasticsearch documents with query: {query}",
- {
- "stage": "fetching_documents",
- "query": query,
- "indices": indices,
- "max_documents": max_docs,
- },
- )
-
- # Index documents
- documents_indexed = 0
- documents_skipped = 0
- skipped_documents = []
-
- try:
- # Connect to Elasticsearch
- if not await elasticsearch_connector.connect():
- await task_logger.log_task_failure(
- log_entry,
- "Failed to connect to Elasticsearch",
- "Connection failed",
- {"error_type": "ConnectionError"},
+ existing = await session.execute(
+ select(Document).where(
+ Document.search_space_id == search_space_id,
+ Document.content_hash == content_hash,
)
- return 0, "Failed to connect to Elasticsearch"
-
- # Search documents
- doc_count = 0
- async for elasticsearch_doc in elasticsearch_connector.search_documents(
- query=query, indices=indices, size=max_docs, fields=fields
- ):
- try:
- if doc_count >= max_docs:
- break
-
- # Convert ElasticsearchDocument to markdown content
- document_content = elasticsearch_doc.content
- document_title = elasticsearch_doc.title
-
- if not document_content.strip():
- logger.warning(
- f"Skipping document with no content: {document_title}"
- )
- skipped_documents.append(f"{document_title} (no content)")
- documents_skipped += 1
- continue
-
- # Generate content hash
- content_hash = generate_content_hash(
- document_content, search_space_id
- )
-
- # Check if document already exists
- existing_document_by_hash = await check_duplicate_document_by_hash(
- session, content_hash
- )
- if existing_document_by_hash:
- logger.info(
- f"Document with content hash {content_hash} already exists for document {document_title}. Skipping processing."
- )
- documents_skipped += 1
- continue
-
- # Generate summary with metadata
- user_llm = await get_user_long_context_llm(session, user_id)
-
- if user_llm:
- # Create metadata for summary generation
- document_metadata = {
- "document_title": document_title,
- "document_type": "Elasticsearch Document",
- "connector_type": "Elasticsearch",
- "source_index": elasticsearch_doc.metadata.get(
- "elasticsearch_index", "unknown"
- ),
- }
- document_metadata.update(elasticsearch_doc.metadata)
-
- (
- summary_content,
- summary_embedding,
- ) = await generate_document_summary(
- document_content, user_llm, document_metadata
- )
- else:
- # Fallback to simple summary if no LLM configured
- summary_content = (
- f"Elasticsearch Document: {document_title}\n\n"
- )
- if elasticsearch_doc.metadata.get("elasticsearch_index"):
- summary_content += f"Index: {elasticsearch_doc.metadata['elasticsearch_index']}\n"
-
- # Add content preview
- content_preview = document_content[:300]
- if len(document_content) > 300:
- content_preview += "..."
- summary_content += f"Content: {content_preview}\n"
-
- summary_embedding = config.embedding_model_instance.embed(
- summary_content
- )
-
- # Create document chunks
- chunks = await create_document_chunks(document_content)
-
- # Create document
- document = Document(
- search_space_id=search_space_id,
- title=f"Elasticsearch Document - {document_title}",
- document_type=DocumentType.ELASTICSEARCH_CONNECTOR,
- document_metadata={
- "document_title": document_title,
- "source_index": elasticsearch_doc.metadata.get(
- "elasticsearch_index", "unknown"
- ),
- "document_id": elasticsearch_doc.metadata.get(
- "elasticsearch_id", "unknown"
- ),
- "indexed_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
- **elasticsearch_doc.metadata,
- },
- content=summary_content,
- content_hash=content_hash,
- embedding=summary_embedding,
- chunks=chunks,
- )
-
- session.add(document)
- documents_indexed += 1
- doc_count += 1
- logger.info(f"Successfully indexed new document {document_title}")
-
- # Add small delay to prevent overwhelming Elasticsearch
- if doc_count % 10 == 0:
- await asyncio.sleep(0.1)
-
- except Exception as e:
- logger.error(
- f"Error processing document {elasticsearch_doc.title if 'elasticsearch_doc' in locals() else 'Unknown'}: {e!s}",
- exc_info=True,
- )
- skipped_documents.append(
- f"{elasticsearch_doc.title if 'elasticsearch_doc' in locals() else 'Unknown'} (processing error)"
- )
- documents_skipped += 1
- continue
-
- logger.info(f"Fetched {doc_count} documents from Elasticsearch")
-
- except Exception as e:
- logger.error(f"Error fetching documents from Elasticsearch: {e}")
- await task_logger.log_task_failure(
- log_entry,
- f"Error fetching documents from Elasticsearch: {e}",
- "Fetch error",
- {"error_type": "FetchError"},
)
- return 0, f"Error fetching documents from Elasticsearch: {e}"
- finally:
- await elasticsearch_connector.disconnect()
+ if existing.scalar_one_or_none():
+ continue
- total_processed = documents_indexed
- if total_processed > 0:
- await update_connector_last_indexed(session, connector, update_last_indexed)
+ title = str(
+ source.get("title")
+ or source.get("name")
+ or hit.get("_id")
+ or "Elasticsearch Document"
+ )
+
+ document_metadata = {
+ "connector_id": connector_id,
+ "connector_name": connector.name,
+ "source": "elasticsearch",
+ "index": hit.get("_index"),
+ "document_id": hit.get("_id"),
+ "score": hit.get("_score"),
+ "search_parameters": {
+ "indices": params.indices,
+ "query": params.query,
+ "search_fields": params.search_fields,
+ "start_date": params.start_date,
+ "end_date": params.end_date,
+ },
+ "ingested_at": ingested_at,
+ }
+
+ new_document = Document(
+ title=title,
+ document_type=DocumentType.ELASTICSEARCH_CONNECTOR,
+ document_metadata=document_metadata,
+ content=content,
+ content_hash=content_hash,
+ search_space_id=search_space_id,
+ )
+ session.add(new_document)
+ await session.flush()
+
+ chunk = Chunk(content=content, document_id=new_document.id)
+ session.add(chunk)
+
+ documents_processed += 1
await session.commit()
-
- await task_logger.log_task_success(
- log_entry,
- f"Successfully completed Elasticsearch indexing for connector {connector_id}",
- {
- "documents_processed": total_processed,
- "documents_indexed": documents_indexed,
- "documents_skipped": documents_skipped,
- "skipped_documents_count": len(skipped_documents),
- },
- )
-
- logger.info(
- f"Elasticsearch indexing completed: {documents_indexed} new documents, {documents_skipped} skipped"
- )
- return total_processed, None
-
- except SQLAlchemyError as db_error:
+ return documents_processed, None
+ except Exception as exc:
await session.rollback()
- await task_logger.log_task_failure(
- log_entry,
- f"Database error during Elasticsearch indexing for connector {connector_id}",
- str(db_error),
- {"error_type": "SQLAlchemyError"},
- )
- logger.error(f"Database error: {db_error!s}", exc_info=True)
- return 0, f"Database error: {db_error!s}"
- except Exception as e:
- await session.rollback()
- await task_logger.log_task_failure(
- log_entry,
- f"Failed to index Elasticsearch documents for connector {connector_id}",
- str(e),
- {"error_type": type(e).__name__},
- )
- logger.error(f"Failed to index Elasticsearch documents: {e!s}", exc_info=True)
- return 0, f"Failed to index Elasticsearch documents: {e!s}"
+ return 0, f"Failed to index Elasticsearch documents: {exc!s}"
diff --git a/surfsense_web/app/dashboard/[search_space_id]/connectors/[connector_id]/edit/page.tsx b/surfsense_web/app/dashboard/[search_space_id]/connectors/[connector_id]/edit/page.tsx
index 87844e8c8..a2dcbab14 100644
--- a/surfsense_web/app/dashboard/[search_space_id]/connectors/[connector_id]/edit/page.tsx
+++ b/surfsense_web/app/dashboard/[search_space_id]/connectors/[connector_id]/edit/page.tsx
@@ -271,6 +271,16 @@ export default function EditConnectorPage() {
placeholder="API Key..."
/>
)}
+ {/* == Elasticsearch == */}
+ {connector.connector_type === "ELASTICSEARCH_CONNECTOR" && (
+
+ )}