Skip to content

Commit a1a9496

Browse files
committed
Fix Langflow DLS ingestion preflight
1 parent 58cc54f commit a1a9496

10 files changed

Lines changed: 198 additions & 34 deletions

Makefile

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -791,6 +791,7 @@ test-ci: ensure-langflow-data ensure-backend-volumes ## Start infra, run integra
791791
GOOGLE_OAUTH_CLIENT_ID="" \
792792
GOOGLE_OAUTH_CLIENT_SECRET="" \
793793
OPENSEARCH_HOST=localhost OPENSEARCH_PORT=9200 \
794+
LANGFLOW_OPENSEARCH_HOST=opensearch LANGFLOW_OPENSEARCH_PORT=9200 \
794795
OPENSEARCH_USERNAME=admin OPENSEARCH_PASSWORD=$${OPENSEARCH_PASSWORD} \
795796
DISABLE_STARTUP_INGEST=$${DISABLE_STARTUP_INGEST:-true} \
796797
uv run pytest tests/integration/core -vv -s -o log_cli=true --log-cli-level=DEBUG; \
@@ -915,6 +916,7 @@ test-ci-local: ensure-langflow-data ensure-backend-volumes ## Same as test-ci bu
915916
GOOGLE_OAUTH_CLIENT_ID="" \
916917
GOOGLE_OAUTH_CLIENT_SECRET="" \
917918
OPENSEARCH_HOST=localhost OPENSEARCH_PORT=9200 \
919+
LANGFLOW_OPENSEARCH_HOST=opensearch LANGFLOW_OPENSEARCH_PORT=9200 \
918920
OPENSEARCH_USERNAME=admin OPENSEARCH_PASSWORD=$${OPENSEARCH_PASSWORD} \
919921
DISABLE_STARTUP_INGEST=$${DISABLE_STARTUP_INGEST:-true} \
920922
uv run pytest tests/integration/core -vv -s -o log_cli=true --log-cli-level=DEBUG; \

flows/components/opensearch_multimodal.py

Lines changed: 16 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,9 +7,6 @@
77
from concurrent.futures import ThreadPoolExecutor, as_completed
88
from typing import Any
99

10-
from opensearchpy import OpenSearch, helpers
11-
from opensearchpy.exceptions import OpenSearchException, RequestError
12-
1310
from lfx.base.vectorstores.model import LCVectorStoreComponent, check_cached_vector_store
1411
from lfx.base.vectorstores.vector_store_connection_decorator import vector_store_connection
1512
from lfx.io import (
@@ -25,6 +22,8 @@
2522
)
2623
from lfx.log import logger
2724
from lfx.schema.data import Data
25+
from opensearchpy import OpenSearch, helpers
26+
from opensearchpy.exceptions import OpenSearchException, RequestError
2827

2928
REQUEST_TIMEOUT = 60
3029
MAX_RETRIES = 5
@@ -813,6 +812,13 @@ def _bulk_ingest_embeddings(
813812
helpers.bulk(client, requests, max_chunk_bytes=max_chunk_bytes)
814813
return return_ids
815814

815+
def _log_index_admin_skip(self, operation: str, error: Exception) -> None:
816+
"""Log index-admin operations that may be blocked under filter-level DLS."""
817+
logger.warning(
818+
f"[OpenSearchMultimodel] Could not run index-admin operation '{operation}': {error}. "
819+
"Assuming the backend pre-created the required index/mapping and continuing."
820+
)
821+
816822
# ---------- param helpers ----------
817823
def _parse_int_param(self, attr_name: str, default: int) -> int:
818824
"""Parse a string attribute to int, returning *default* on failure."""
@@ -1228,8 +1234,14 @@ def _embed(text: str) -> list[float]:
12281234
)
12291235

12301236
# Ensure index exists with baseline mapping (index.knn: true is required for vector search)
1237+
index_exists = True
1238+
try:
1239+
index_exists = bool(client.indices.exists(index=self.index_name))
1240+
except OpenSearchException as exists_error:
1241+
self._log_index_admin_skip("indices.exists", exists_error)
1242+
12311243
try:
1232-
if not client.indices.exists(index=self.index_name):
1244+
if not index_exists:
12331245
self.log(f"Creating index '{self.index_name}' with base mapping")
12341246
client.indices.create(index=self.index_name, body=mapping)
12351247
except RequestError as creation_error:

flows/ingestion_flow.json

Lines changed: 1 addition & 1 deletion
Large diffs are not rendered by default.

flows/openrag_agent.json

Lines changed: 1 addition & 1 deletion
Large diffs are not rendered by default.

flows/openrag_nudges.json

Lines changed: 1 addition & 1 deletion
Large diffs are not rendered by default.

flows/openrag_url_mcp.json

Lines changed: 1 addition & 1 deletion
Large diffs are not rendered by default.

src/services/langflow_file_service.py

Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ def __init__(self, flows_service=None, docling_service=None):
2222
self.flows_service = flows_service
2323
self.docling_service = docling_service
2424
self.flow_id_url_ingest = LANGFLOW_URL_INGEST_FLOW_ID
25+
self._embedding_dimension_cache: dict[str, int] = {}
2526

2627
_TRANSIENT_STATUS_CODES = {408, 429, 500, 502, 503, 504}
2728

@@ -71,6 +72,77 @@ def merge_ui_ingest_settings_into_tweaks(
7172

7273
return final_tweaks
7374

75+
async def _detect_embedding_dimensions(
76+
self,
77+
embedding_model: str,
78+
embedding_provider: str | None,
79+
) -> int:
80+
"""Generate one probe embedding so mapping dimensions match the provider."""
81+
from services.models_service import ModelsService
82+
83+
cache_key = f"{embedding_provider or ''}:{embedding_model}"
84+
cached = self._embedding_dimension_cache.get(cache_key)
85+
if cached:
86+
return cached
87+
88+
litellm_model_name = await ModelsService().get_litellm_model_name(
89+
embedding_model,
90+
provider=embedding_provider,
91+
)
92+
response = await clients.patched_embedding_client.embeddings.create(
93+
model=litellm_model_name,
94+
input=["dimension probe"],
95+
)
96+
if not response.data:
97+
raise RuntimeError("Embedding provider returned no data for dimension probe")
98+
99+
first = response.data[0]
100+
embedding = first["embedding"] if isinstance(first, dict) else first.embedding
101+
dimensions = len(embedding)
102+
if dimensions <= 0:
103+
raise RuntimeError("Embedding provider returned an empty dimension probe")
104+
105+
self._embedding_dimension_cache[cache_key] = dimensions
106+
return dimensions
107+
108+
async def _ensure_langflow_ingest_index(self, embedding_model: str | None) -> None:
109+
"""Pre-create index mappings Langflow cannot manage with a DLS JWT."""
110+
if clients.opensearch is None:
111+
logger.debug("[LF] OpenSearch admin client unavailable; skipping ingest preflight")
112+
return
113+
114+
try:
115+
from config.embedding_constants import OPENAI_DEFAULT_EMBEDDING_MODEL
116+
from config.settings import get_index_name, get_openrag_config
117+
from utils.embedding_fields import ensure_embedding_field_exists
118+
from utils.embeddings import create_index_body
119+
120+
config = get_openrag_config()
121+
index_name = get_index_name()
122+
model_name = embedding_model or OPENAI_DEFAULT_EMBEDDING_MODEL
123+
embedding_dimensions = await self._detect_embedding_dimensions(
124+
model_name,
125+
config.knowledge.embedding_provider,
126+
)
127+
if not await clients.opensearch.indices.exists(index=index_name):
128+
await clients.opensearch.indices.create(
129+
index=index_name,
130+
body=await create_index_body(model_name, embedding_dimensions),
131+
)
132+
133+
await ensure_embedding_field_exists(
134+
clients.opensearch,
135+
model_name,
136+
index_name,
137+
embedding_dimensions,
138+
)
139+
except Exception as e:
140+
logger.warning(
141+
"[LF] Failed to preconfigure OpenSearch index before Langflow ingest",
142+
embedding_model=embedding_model,
143+
error=str(e),
144+
)
145+
74146
async def upload_user_file(self, file_tuple, jwt_token: str | None = None) -> dict[str, Any]:
75147
"""Upload a file using Langflow Files API v2: POST /api/v2/files.
76148
Returns JSON with keys: id, name, path, size, provider.
@@ -229,6 +301,7 @@ async def run_ingestion_flow(
229301
await add_provider_credentials_to_headers(
230302
headers, config, flows_service=self.flows_service, jwt_token=jwt_token
231303
)
304+
await self._ensure_langflow_ingest_index(embedding_model)
232305
start_time = time.time()
233306
logger.info(
234307
"[INGEST] Run started",
@@ -353,6 +426,7 @@ async def run_url_ingestion_flow(
353426
await add_provider_credentials_to_headers(
354427
headers, config, flows_service=self.flows_service, jwt_token=jwt_token
355428
)
429+
await self._ensure_langflow_ingest_index(embedding_model)
356430

357431
logger.info(
358432
"[LF] Running URL ingestion flow",

src/utils/embeddings.py

Lines changed: 44 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -1,37 +1,58 @@
1+
from utils.embedding_fields import build_knn_vector_field, get_embedding_field_name
12
from utils.logging_config import get_logger
23

34
logger = get_logger(__name__)
45

56

6-
async def create_index_body() -> dict:
7+
async def create_index_body(
8+
embedding_model: str | None = None,
9+
embedding_dimensions: int | None = None,
10+
) -> dict:
711
"""Create a static index body configuration.
812
913
Returns:
1014
OpenSearch index body configuration
1115
"""
16+
from config.embedding_constants import OPENAI_DEFAULT_EMBEDDING_MODEL
17+
from config.settings import VECTOR_DIM, get_openrag_config
18+
19+
resolved_embedding_model = (
20+
embedding_model
21+
or get_openrag_config().knowledge.embedding_model
22+
or OPENAI_DEFAULT_EMBEDDING_MODEL
23+
)
24+
25+
properties = {
26+
"document_id": {"type": "keyword"},
27+
"filename": {"type": "keyword"},
28+
"mimetype": {"type": "keyword"},
29+
"page": {"type": "integer"},
30+
"text": {"type": "text"},
31+
# Legacy field - kept for backward compatibility and for clusters where
32+
# Langflow cannot perform mapping updates with a DLS-filtered JWT.
33+
"chunk_embedding": build_knn_vector_field(VECTOR_DIM),
34+
# Track which embedding model was used for this chunk
35+
"embedding_model": {"type": "keyword"},
36+
"embedding_dimensions": {"type": "integer"},
37+
"source_url": {"type": "keyword"},
38+
"connector_type": {"type": "keyword"},
39+
"owner": {"type": "keyword"},
40+
"owner_email": {"type": "keyword"},
41+
"allowed_users": {"type": "keyword"},
42+
"allowed_groups": {"type": "keyword"},
43+
"allowed_principals": {"type": "keyword"},
44+
"created_time": {"type": "date"},
45+
"modified_time": {"type": "date"},
46+
"indexed_time": {"type": "date"},
47+
"metadata": {"type": "object"},
48+
}
49+
50+
if embedding_dimensions:
51+
properties[get_embedding_field_name(resolved_embedding_model)] = build_knn_vector_field(
52+
embedding_dimensions
53+
)
1254

1355
return {
1456
"settings": {"index": {"knn": True}, "number_of_shards": 1, "number_of_replicas": 0},
15-
"mappings": {
16-
"properties": {
17-
"document_id": {"type": "keyword"},
18-
"filename": {"type": "keyword"},
19-
"mimetype": {"type": "keyword"},
20-
"page": {"type": "integer"},
21-
"text": {"type": "text"},
22-
# Track which embedding model was used for this chunk
23-
"embedding_model": {"type": "keyword"},
24-
"embedding_dimensions": {"type": "integer"},
25-
"source_url": {"type": "keyword"},
26-
"connector_type": {"type": "keyword"},
27-
"owner": {"type": "keyword"},
28-
"allowed_users": {"type": "keyword"},
29-
"allowed_groups": {"type": "keyword"},
30-
"allowed_principals": {"type": "keyword"},
31-
"created_time": {"type": "date"},
32-
"modified_time": {"type": "date"},
33-
"indexed_time": {"type": "date"},
34-
"metadata": {"type": "object"},
35-
}
36-
},
57+
"mappings": {"properties": properties},
3758
}

tests/unit/test_embedding_fields.py

Lines changed: 26 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,11 +8,12 @@
88
JVector/DiskANN method configuration with only the dimension varying per
99
embedding model.
1010
"""
11-
from typing import Any, Dict
11+
from types import SimpleNamespace
12+
from typing import Any
1213

1314
import pytest
1415

15-
from utils.embedding_fields import build_knn_vector_field
16+
from utils.embedding_fields import build_knn_vector_field, get_embedding_field_name
1617

1718

1819
class TestBuildKnnVectorFieldStructure:
@@ -99,6 +100,28 @@ class TestBuildKnnVectorFieldCallSitesMatch:
99100
def test_index_body_uses_helper_output(self) -> None:
100101
from config.settings import INDEX_BODY, VECTOR_DIM
101102

102-
chunk_field: Dict[str, Any] = INDEX_BODY["mappings"]["properties"]["chunk_embedding"]
103+
chunk_field: dict[str, Any] = INDEX_BODY["mappings"]["properties"][
104+
"chunk_embedding"
105+
]
103106
expected = build_knn_vector_field(VECTOR_DIM)
104107
assert chunk_field == expected
108+
109+
@pytest.mark.asyncio
110+
async def test_create_index_body_precreates_configured_embedding_field(
111+
self, monkeypatch: pytest.MonkeyPatch
112+
) -> None:
113+
monkeypatch.setattr(
114+
"config.settings.get_openrag_config",
115+
lambda: SimpleNamespace(
116+
knowledge=SimpleNamespace(embedding_model="text-embedding-3-large")
117+
),
118+
)
119+
120+
from utils.embeddings import create_index_body
121+
122+
body = await create_index_body("text-embedding-3-large", 3072)
123+
properties = body["mappings"]["properties"]
124+
embedding_field = get_embedding_field_name("text-embedding-3-large")
125+
126+
assert properties[embedding_field] == build_knn_vector_field(3072)
127+
assert properties["owner_email"] == {"type": "keyword"}

tests/unit/test_langflow_file_service_two_phase.py

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
Docling fails / expires / times out.
77
"""
88

9+
from types import SimpleNamespace
910
from unittest.mock import AsyncMock
1011

1112
import pytest
@@ -231,6 +232,37 @@ def test_processor_accepts_injected_polling_service():
231232
assert processor.docling_polling_service is injected
232233

233234

235+
@pytest.mark.asyncio
236+
async def test_langflow_preflight_detects_embedding_dimensions_with_probe(monkeypatch):
237+
class FakeEmbeddings:
238+
def __init__(self):
239+
self.calls = []
240+
241+
async def create(self, model, input):
242+
self.calls.append((model, input))
243+
return SimpleNamespace(data=[SimpleNamespace(embedding=[0.1] * 7)])
244+
245+
fake_embeddings = FakeEmbeddings()
246+
fake_client = SimpleNamespace(embeddings=fake_embeddings)
247+
248+
async def fake_get_litellm_model_name(self, model_name, provider=None, strict=False):
249+
assert model_name == "provider/model"
250+
assert provider == "provider"
251+
return "provider/provider/model"
252+
253+
monkeypatch.setattr("config.settings.clients._patched_async_client", fake_client)
254+
monkeypatch.setattr(
255+
"services.models_service.ModelsService.get_litellm_model_name",
256+
fake_get_litellm_model_name,
257+
)
258+
259+
svc = LangflowFileService(docling_service=AsyncMock())
260+
261+
assert await svc._detect_embedding_dimensions("provider/model", "provider") == 7
262+
assert await svc._detect_embedding_dimensions("provider/model", "provider") == 7
263+
assert fake_embeddings.calls == [("provider/provider/model", ["dimension probe"])]
264+
265+
234266
@pytest.mark.asyncio
235267
async def test_task_service_threads_polling_service_to_processor(monkeypatch):
236268
"""TaskService.create_langflow_upload_task must forward its injected

0 commit comments

Comments
 (0)