validops-east-1 commited on
Commit
6fa7b33
·
1 Parent(s): cd6f706

chore: remaining working-tree changes

Browse files
app/api/v1/vector_stores.py CHANGED
@@ -14,7 +14,11 @@ from fastapi import (
14
  status,
15
  )
16
 
17
- from app.api.deps import get_vector_store_service
 
 
 
 
18
  from app.core.logger import get_logger
19
  from app.models.domain import ConversionError
20
  from app.models.schemas import (
@@ -38,18 +42,22 @@ router = APIRouter()
38
  logger = get_logger(__name__)
39
 
40
 
41
- async def _process_pdf_bytes(raw: bytes, source: str, clean_content: bool) -> str | ConversionError:
 
 
 
 
 
 
42
  if len(raw) < 5 or raw[:5] != b"%PDF-":
43
  return ConversionError(source=source, error_type="ValueError", message="Not a valid PDF", duration_ms=0)
44
  loop = asyncio.get_running_loop()
45
- converter = ConverterService()
46
- result = await loop.run_in_executor(None, converter.convert_stream, raw, source)
47
  if isinstance(result, ConversionError):
48
  return result
49
  text = result.markdown
50
  if clean_content:
51
- text_cleaner = TextCleanerService()
52
- text = await loop.run_in_executor(None, text_cleaner.clean, text)
53
  return text
54
 
55
 
@@ -214,6 +222,8 @@ async def ingest_pdf_document(
214
  chunk_overlap: int = Form(64, ge=0, le=512),
215
  clean_content: bool = Query(True, description="Clean markdown text after PDF conversion"),
216
  vector_store_service: VectorStoreService = Depends(get_vector_store_service),
 
 
217
  ) -> DocumentIngestResponse:
218
  record = vector_store_service.get_store(store_id)
219
  if record is None:
@@ -229,7 +239,9 @@ async def ingest_pdf_document(
229
  finally:
230
  await file.close()
231
 
232
- text = await _process_pdf_bytes(raw, file.filename, clean_content)
 
 
233
  if isinstance(text, ConversionError):
234
  return DocumentIngestResponse(
235
  success=False,
@@ -279,6 +291,8 @@ async def ingest_pdf_url(
279
  body: DocumentIngestUrlRequest,
280
  clean_content: bool = Query(True, description="Clean markdown text after PDF conversion"),
281
  vector_store_service: VectorStoreService = Depends(get_vector_store_service),
 
 
282
  ) -> DocumentIngestResponse:
283
  record = vector_store_service.get_store(store_id)
284
  if record is None:
@@ -296,7 +310,9 @@ async def ingest_pdf_url(
296
  error=f"Failed to fetch PDF from URL: {exc}",
297
  )
298
 
299
- text = await _process_pdf_bytes(raw, body.url, clean_content)
 
 
300
  if isinstance(text, ConversionError):
301
  return DocumentIngestResponse(
302
  success=False,
 
14
  status,
15
  )
16
 
17
+ from app.api.deps import (
18
+ get_converter_service,
19
+ get_text_cleaner_service,
20
+ get_vector_store_service,
21
+ )
22
  from app.core.logger import get_logger
23
  from app.models.domain import ConversionError
24
  from app.models.schemas import (
 
42
  logger = get_logger(__name__)
43
 
44
 
45
+ async def _process_pdf_bytes(
46
+ raw: bytes,
47
+ source: str,
48
+ clean_content: bool,
49
+ converter_service: ConverterService,
50
+ text_cleaner_service: TextCleanerService,
51
+ ) -> str | ConversionError:
52
  if len(raw) < 5 or raw[:5] != b"%PDF-":
53
  return ConversionError(source=source, error_type="ValueError", message="Not a valid PDF", duration_ms=0)
54
  loop = asyncio.get_running_loop()
55
+ result = await loop.run_in_executor(None, converter_service.convert_stream, raw, source)
 
56
  if isinstance(result, ConversionError):
57
  return result
58
  text = result.markdown
59
  if clean_content:
60
+ text = await loop.run_in_executor(None, text_cleaner_service.clean, text)
 
61
  return text
62
 
63
 
 
222
  chunk_overlap: int = Form(64, ge=0, le=512),
223
  clean_content: bool = Query(True, description="Clean markdown text after PDF conversion"),
224
  vector_store_service: VectorStoreService = Depends(get_vector_store_service),
225
+ converter_service: ConverterService = Depends(get_converter_service),
226
+ text_cleaner_service: TextCleanerService = Depends(get_text_cleaner_service),
227
  ) -> DocumentIngestResponse:
228
  record = vector_store_service.get_store(store_id)
229
  if record is None:
 
239
  finally:
240
  await file.close()
241
 
242
+ text = await _process_pdf_bytes(
243
+ raw, file.filename, clean_content, converter_service, text_cleaner_service
244
+ )
245
  if isinstance(text, ConversionError):
246
  return DocumentIngestResponse(
247
  success=False,
 
291
  body: DocumentIngestUrlRequest,
292
  clean_content: bool = Query(True, description="Clean markdown text after PDF conversion"),
293
  vector_store_service: VectorStoreService = Depends(get_vector_store_service),
294
+ converter_service: ConverterService = Depends(get_converter_service),
295
+ text_cleaner_service: TextCleanerService = Depends(get_text_cleaner_service),
296
  ) -> DocumentIngestResponse:
297
  record = vector_store_service.get_store(store_id)
298
  if record is None:
 
310
  error=f"Failed to fetch PDF from URL: {exc}",
311
  )
312
 
313
+ text = await _process_pdf_bytes(
314
+ raw, body.url, clean_content, converter_service, text_cleaner_service
315
+ )
316
  if isinstance(text, ConversionError):
317
  return DocumentIngestResponse(
318
  success=False,
app/api/v1/web_search.py CHANGED
@@ -20,9 +20,11 @@ from app.services.web_search_service import WebSearchService
20
  router = APIRouter()
21
  _logger = get_logger(__name__)
22
 
 
 
23
 
24
  def get_search_service() -> WebSearchService:
25
- return WebSearchService()
26
 
27
 
28
  @router.post(
 
20
  router = APIRouter()
21
  _logger = get_logger(__name__)
22
 
23
+ _search_service = WebSearchService()
24
+
25
 
26
  def get_search_service() -> WebSearchService:
27
+ return _search_service
28
 
29
 
30
  @router.post(
app/services/semantic_router_service.py CHANGED
@@ -1,10 +1,13 @@
1
  from __future__ import annotations
2
 
3
  import logging
4
- from collections import defaultdict
5
- from typing import Any
6
 
7
  import numpy as np
 
 
 
 
8
 
9
  from app.services.embeddings_service import EmbeddingService
10
 
@@ -13,6 +16,33 @@ _logger = logging.getLogger(__name__)
13
  _DIMENSION = 384
14
 
15
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
16
  class SemanticRouterService:
17
  def __init__(self, embedding_service: EmbeddingService) -> None:
18
  self._embedding_service = embedding_service
@@ -27,70 +57,76 @@ class SemanticRouterService:
27
  return {"success": False, "name": None, "models": [], "error": "Query is empty."}
28
  if not routes:
29
  return {"success": False, "name": None, "models": [], "error": "No routes provided."}
30
-
31
  if not self._embedding_service.is_loaded(_DIMENSION):
32
  return {"success": False, "name": None, "models": [], "error": "Embedding model not loaded."}
33
 
34
- all_utterances: list[str] = []
35
- utterance_to_route: list[int] = []
 
36
 
37
- for idx, route in enumerate(routes):
38
- utterances = route.get("utterances", [])
39
- if not utterances:
40
- continue
41
- all_utterances.extend(utterances)
42
- utterance_to_route.extend([idx] * len(utterances))
43
 
44
- if not all_utterances:
45
- return {"success": False, "name": None, "models": [], "error": "No utterances found in any route."}
 
46
 
47
  try:
48
- all_texts = all_utterances + [query]
49
- embeddings = self._embedding_service.generate_embedding(all_texts, _DIMENSION)
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
50
  except Exception as exc:
51
- _logger.error("Embedding failed: %s", exc)
52
  return {"success": False, "name": None, "models": [], "error": str(exc)}
53
 
54
- utterance_embs = np.array(embeddings[:-1], dtype=np.float32)
55
- query_emb = np.array(embeddings[-1], dtype=np.float32).reshape(1, -1)
56
-
57
- similarities = (utterance_embs @ query_emb.T).flatten()
58
- best_idx = int(np.argmax(similarities))
59
- best_score = float(similarities[best_idx])
60
 
61
- # Route-level scoring: take the max score per route
62
- route_scores: dict[int, list[float]] = defaultdict(list)
63
- for i, score in enumerate(similarities):
64
- route_scores[utterance_to_route[i]].append(float(score))
65
-
66
- route_best_scores = {rid: max(scores) for rid, scores in route_scores.items()}
67
- sorted_routes = sorted(route_best_scores.items(), key=lambda x: x[1], reverse=True)
68
- best_route_score = sorted_routes[0][1]
69
- second_best_route_score = sorted_routes[1][1] if len(sorted_routes) > 1 else 1.0
70
- route_margin = best_route_score - second_best_route_score
71
-
72
- if best_score < threshold or route_margin < 0.001:
73
  return {
74
  "success": True,
75
  "name": None,
76
  "models": [],
77
  "error": None,
78
- "confidence": best_score,
79
- "margin": route_margin,
80
  "threshold": threshold,
81
- "matched_utterance": all_utterances[best_idx],
82
  }
83
 
84
- matched_route_idx = utterance_to_route[best_idx]
85
- matched_route = routes[matched_route_idx]
86
-
87
  return {
88
  "success": True,
89
- "name": matched_route.get("name"),
90
- "models": matched_route.get("models", []),
91
  "error": None,
92
- "confidence": best_score,
93
- "margin": route_margin,
94
  "threshold": threshold,
95
- "matched_utterance": all_utterances[best_idx],
96
  }
 
1
  from __future__ import annotations
2
 
3
  import logging
4
+ from typing import Any, Dict, List, Optional
 
5
 
6
  import numpy as np
7
+ from semantic_router import Route
8
+ from semantic_router.encoders import DenseEncoder
9
+ from semantic_router.linear import similarity_matrix
10
+ from semantic_router.routers import SemanticRouter
11
 
12
  from app.services.embeddings_service import EmbeddingService
13
 
 
16
  _DIMENSION = 384
17
 
18
 
19
+ class _EmbeddingServiceEncoder(DenseEncoder):
20
+ """Adapter exposing the project's EmbeddingService as a semantic-router encoder.
21
+
22
+ Wrapping the already-loaded sentence-transformers model avoids loading a second
23
+ embedding model just for routing.
24
+ """
25
+
26
+ def __init__(
27
+ self,
28
+ embedding_service: EmbeddingService,
29
+ dimension: int = _DIMENSION,
30
+ score_threshold: Optional[float] = None,
31
+ ) -> None:
32
+ super().__init__(
33
+ name="granite-embedding-small-english-r2",
34
+ score_threshold=score_threshold,
35
+ )
36
+ self._embedding_service = embedding_service
37
+ self._dimension = dimension
38
+
39
+ def __call__(self, docs: List[Any]) -> List[List[float]]:
40
+ return self._embedding_service.generate_embedding([str(d) for d in docs], self._dimension)
41
+
42
+ async def acall(self, docs: List[Any]) -> List[List[float]]:
43
+ return self(docs)
44
+
45
+
46
  class SemanticRouterService:
47
  def __init__(self, embedding_service: EmbeddingService) -> None:
48
  self._embedding_service = embedding_service
 
57
  return {"success": False, "name": None, "models": [], "error": "Query is empty."}
58
  if not routes:
59
  return {"success": False, "name": None, "models": [], "error": "No routes provided."}
 
60
  if not self._embedding_service.is_loaded(_DIMENSION):
61
  return {"success": False, "name": None, "models": [], "error": "Embedding model not loaded."}
62
 
63
+ valid_routes = [r for r in routes if r.get("utterances")]
64
+ if not valid_routes:
65
+ return {"success": False, "name": None, "models": [], "error": "No utterances found in any route."}
66
 
67
+ name_to_models: Dict[str, List[str]] = {r["name"]: r.get("models", []) for r in routes}
 
 
 
 
 
68
 
69
+ encoder = _EmbeddingServiceEncoder(self._embedding_service, score_threshold=threshold)
70
+ lib_routes = [Route(name=r["name"], utterances=r["utterances"]) for r in valid_routes]
71
+ total_utterances = sum(len(r.utterances) for r in lib_routes)
72
 
73
  try:
74
+ router = SemanticRouter(
75
+ encoder=encoder,
76
+ routes=lib_routes,
77
+ auto_sync="local",
78
+ aggregation="max",
79
+ top_k=total_utterances,
80
+ )
81
+ # Encode the query once, then reuse the vector for both the router
82
+ # decision and the per-route margin / matched-utterance computation.
83
+ query_vector = np.asarray(encoder([query])[0], dtype=np.float32)
84
+ choice = router(vector=query_vector)
85
+ if isinstance(choice, list):
86
+ choice = choice[0] if choice else None
87
+
88
+ sims = similarity_matrix(query_vector, router.index.index)
89
+ scores_by_route: Dict[str, float] = {}
90
+ for route_name, score in zip(router.index.routes, sims):
91
+ scores_by_route[str(route_name)] = max(
92
+ scores_by_route.get(str(route_name), float("-inf")), float(score)
93
+ )
94
+ sorted_routes = sorted(scores_by_route.items(), key=lambda x: x[1], reverse=True)
95
+ best_score = sorted_routes[0][1]
96
+ second_best = sorted_routes[1][1] if len(sorted_routes) > 1 else 1.0
97
+ margin = best_score - second_best
98
+ best_utt_idx = int(np.argmax(sims))
99
+ matched_utterance = str(router.index.utterances[best_utt_idx])
100
  except Exception as exc:
101
+ _logger.error("Semantic routing failed: %s", exc)
102
  return {"success": False, "name": None, "models": [], "error": str(exc)}
103
 
104
+ name = choice.name if choice is not None else None
105
+ confidence = (
106
+ float(choice.similarity_score)
107
+ if choice is not None and choice.similarity_score is not None
108
+ else best_score
109
+ )
110
 
111
+ if name is None or margin < 0.001:
 
 
 
 
 
 
 
 
 
 
 
112
  return {
113
  "success": True,
114
  "name": None,
115
  "models": [],
116
  "error": None,
117
+ "confidence": confidence,
118
+ "margin": margin,
119
  "threshold": threshold,
120
+ "matched_utterance": matched_utterance,
121
  }
122
 
 
 
 
123
  return {
124
  "success": True,
125
+ "name": name,
126
+ "models": name_to_models.get(name, []),
127
  "error": None,
128
+ "confidence": confidence,
129
+ "margin": margin,
130
  "threshold": threshold,
131
+ "matched_utterance": matched_utterance,
132
  }
pyproject.toml CHANGED
@@ -24,6 +24,7 @@ dependencies = [
24
  "pyfiglet>=1.0.0",
25
  "rich>=13.7",
26
  "numpy>=1.26.0",
 
27
  "rapidocr-onnxruntime>=1.4.4",
28
  "onnxruntime>=1.18.0",
29
  "pillow>=10.0.0",
 
24
  "pyfiglet>=1.0.0",
25
  "rich>=13.7",
26
  "numpy>=1.26.0",
27
+ "semantic-router>=0.1.16",
28
  "rapidocr-onnxruntime>=1.4.4",
29
  "onnxruntime>=1.18.0",
30
  "pillow>=10.0.0",