From 95f41311d581bc7d286f72e5b21b1e5382a3197c Mon Sep 17 00:00:00 2001 From: Shantanav Mukherjee Date: Fri, 2 Oct 2026 00:13:46 +0530 Subject: [PATCH 1/8] Deduplicate RAG chunks across embedding batches --- python/app/rag/embeddings/vector_db.py | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/python/app/rag/embeddings/vector_db.py b/python/app/rag/embeddings/vector_db.py index 5589baa..969e089 100755 --- a/python/app/rag/embeddings/vector_db.py +++ b/python/app/rag/embeddings/vector_db.py @@ -20,8 +20,11 @@ def __init__(self, collection_name): ) def add_documents(self, documents, batch_size=100): - for i in range(0, len(documents), batch_size): - batch = deduplicate_documents(documents[i : i + batch_size]) + # Deduplicate before batching so repeated chunks across batch boundaries + # do not trigger a second embedding/vector write. + unique_documents = deduplicate_documents(documents) + for i in range(0, len(unique_documents), batch_size): + batch = unique_documents[i : i + batch_size] if not batch: continue From 6d78d4171b9654ad924d85c2fabd19c8a9dec1a0 Mon Sep 17 00:00:00 2001 From: Shantanav Mukherjee Date: Fri, 2 Oct 2026 00:14:10 +0530 Subject: [PATCH 2/8] Test cross-batch RAG chunk deduplication --- python/tests/test_vector_db.py | 40 ++++++++++++++++++++++++++++++++++ 1 file changed, 40 insertions(+) create mode 100644 python/tests/test_vector_db.py diff --git a/python/tests/test_vector_db.py b/python/tests/test_vector_db.py new file mode 100644 index 0000000..949fba0 --- /dev/null +++ b/python/tests/test_vector_db.py @@ -0,0 +1,40 @@ +import unittest +from unittest.mock import MagicMock + +from langchain_core.documents import Document + +from app.rag.embeddings.vector_db import vector_db + + +class VectorDatabaseIngestionTests(unittest.TestCase): + def test_duplicate_chunks_across_batches_are_submitted_once(self): + unique = [ + Document( + page_content=f"Source passage {index}", + metadata={"job_id": "job-1", "chunk_index": index}, + ) + for index in range(100) + ] + duplicated = [ + Document( + page_content=f" source passage {index} ", + metadata={"job_id": "job-1", "chunk_index": index + 100}, + ) + for index in range(100) + ] + collection = MagicMock() + database = vector_db.__new__(vector_db) + database.collection = collection + + database.add_documents(unique + duplicated, batch_size=100) + + submitted = sum( + len(call.args[0]) + for call in collection.add_documents.call_args_list + ) + self.assertEqual(submitted, 100) + self.assertEqual(collection.add_documents.call_count, 1) + + +if __name__ == "__main__": + unittest.main() From 5548e416072d953ee18ac1c9ac1a501f42251d79 Mon Sep 17 00:00:00 2001 From: Shantanav Mukherjee Date: Fri, 2 Oct 2026 00:14:55 +0530 Subject: [PATCH 3/8] Support streaming deduplication across RAG batches --- python/app/rag/embeddings/reranker.py | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/python/app/rag/embeddings/reranker.py b/python/app/rag/embeddings/reranker.py index 552b65d..55a1c38 100755 --- a/python/app/rag/embeddings/reranker.py +++ b/python/app/rag/embeddings/reranker.py @@ -19,14 +19,17 @@ def document_signature(document: Document) -> str: return sha1(normalized_content.encode("utf-8")).hexdigest() -def deduplicate_documents(documents: list[Document]) -> list[Document]: +def deduplicate_documents( + documents: list[Document], + seen_signatures: set[str] | None = None, +) -> list[Document]: """ Remove duplicate documents based on content similarity. - Uses SHA1 hashing of normalized content to detect duplicates. - Safe for low-memory environments (no ML models). + Uses SHA1 hashing of normalized content to detect duplicates. Callers can + pass a shared signature set to deduplicate a stream across bounded batches. """ - seen_signatures: set[str] = set() + seen = seen_signatures if seen_signatures is not None else set() deduplicated: list[Document] = [] for document in documents: @@ -34,10 +37,10 @@ def deduplicate_documents(documents: list[Document]) -> list[Document]: continue signature = document_signature(document) - if signature in seen_signatures: + if signature in seen: continue - seen_signatures.add(signature) + seen.add(signature) deduplicated.append(document) return deduplicated From 415bc31212081b605b950aaea5d0e21816506bc0 Mon Sep 17 00:00:00 2001 From: Shantanav Mukherjee Date: Fri, 2 Oct 2026 00:15:03 +0530 Subject: [PATCH 4/8] Deduplicate RAG chunks incrementally across batches --- python/app/rag/embeddings/vector_db.py | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/python/app/rag/embeddings/vector_db.py b/python/app/rag/embeddings/vector_db.py index 969e089..b303516 100755 --- a/python/app/rag/embeddings/vector_db.py +++ b/python/app/rag/embeddings/vector_db.py @@ -20,11 +20,14 @@ def __init__(self, collection_name): ) def add_documents(self, documents, batch_size=100): - # Deduplicate before batching so repeated chunks across batch boundaries - # do not trigger a second embedding/vector write. - unique_documents = deduplicate_documents(documents) - for i in range(0, len(unique_documents), batch_size): - batch = unique_documents[i : i + batch_size] + # Keep one content-signature set across batches without copying the + # complete document corpus into a second list. + seen_signatures: set[str] = set() + for i in range(0, len(documents), batch_size): + batch = deduplicate_documents( + documents[i : i + batch_size], + seen_signatures=seen_signatures, + ) if not batch: continue From f2b9daac2e7bac27368952cab1eaff4c15d6c593 Mon Sep 17 00:00:00 2001 From: Shantanav Mukherjee Date: Fri, 2 Oct 2026 00:42:21 +0530 Subject: [PATCH 5/8] Use compact bounded signatures for cross-batch RAG deduplication --- python/app/rag/embeddings/reranker.py | 33 ++++++++++++++++++--------- 1 file changed, 22 insertions(+), 11 deletions(-) diff --git a/python/app/rag/embeddings/reranker.py b/python/app/rag/embeddings/reranker.py index 55a1c38..773bc5e 100755 --- a/python/app/rag/embeddings/reranker.py +++ b/python/app/rag/embeddings/reranker.py @@ -5,7 +5,7 @@ from langchain_core.documents import Document -_WHITESPACE_RE = re.compile(r"\s+") +_WHITESPACE_RE = re.compile(r"\\s+") def normalize_text(text: str) -> str: @@ -13,23 +13,24 @@ def normalize_text(text: str) -> str: return _WHITESPACE_RE.sub(" ", text).strip().lower() -def document_signature(document: Document) -> str: - """Generate a signature for a document based on its content.""" +def document_signature(document: Document) -> bytes: + """Generate a compact signature for a document's normalized content.""" normalized_content = normalize_text(document.page_content) - return sha1(normalized_content.encode("utf-8")).hexdigest() + return sha1(normalized_content.encode("utf-8")).digest() def deduplicate_documents( documents: list[Document], - seen_signatures: set[str] | None = None, + seen_signatures: set[bytes] | None = None, + max_seen_signatures: int | None = None, ) -> list[Document]: """ - Remove duplicate documents based on content similarity. + Remove duplicate documents by normalized content. - Uses SHA1 hashing of normalized content to detect duplicates. Callers can - pass a shared signature set to deduplicate a stream across bounded batches. + A shared signature set deduplicates across batches. Its optional cap bounds + extra memory; duplicates remain deduplicated within each individual batch. """ - seen = seen_signatures if seen_signatures is not None else set() + batch_signatures: set[bytes] = set() deduplicated: list[Document] = [] for document in documents: @@ -37,10 +38,20 @@ def deduplicate_documents( continue signature = document_signature(document) - if signature in seen: + if signature in batch_signatures: continue + batch_signatures.add(signature) - seen.add(signature) + if seen_signatures is not None and signature in seen_signatures: + continue + if ( + seen_signatures is not None + and ( + max_seen_signatures is None + or len(seen_signatures) < max_seen_signatures + ) + ): + seen_signatures.add(signature) deduplicated.append(document) return deduplicated From 00a520790ad350ff544c572423ba3aee0cedc788 Mon Sep 17 00:00:00 2001 From: Shantanav Mukherjee Date: Fri, 2 Oct 2026 00:42:29 +0530 Subject: [PATCH 6/8] Bound RAG dedup metadata memory --- python/app/rag/embeddings/vector_db.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/python/app/rag/embeddings/vector_db.py b/python/app/rag/embeddings/vector_db.py index b303516..2c3330e 100755 --- a/python/app/rag/embeddings/vector_db.py +++ b/python/app/rag/embeddings/vector_db.py @@ -22,11 +22,13 @@ def __init__(self, collection_name): def add_documents(self, documents, batch_size=100): # Keep one content-signature set across batches without copying the # complete document corpus into a second list. - seen_signatures: set[str] = set() + max_seen_signatures = int(os.getenv("RAG_DEDUP_MAX_SIGNATURES", "50000")) + seen_signatures: set[bytes] = set() for i in range(0, len(documents), batch_size): batch = deduplicate_documents( documents[i : i + batch_size], seen_signatures=seen_signatures, + max_seen_signatures=max_seen_signatures, ) if not batch: From 7d9435b70c47eb5dbf3c28fc8717ed2a3400f8bd Mon Sep 17 00:00:00 2001 From: Shantanav Mukherjee Date: Fri, 2 Oct 2026 00:42:39 +0530 Subject: [PATCH 7/8] Test the RAG dedup memory cap --- python/tests/test_vector_db.py | 21 ++++++++++++++++++++- 1 file changed, 20 insertions(+), 1 deletion(-) diff --git a/python/tests/test_vector_db.py b/python/tests/test_vector_db.py index 949fba0..c562346 100644 --- a/python/tests/test_vector_db.py +++ b/python/tests/test_vector_db.py @@ -1,5 +1,6 @@ +import os import unittest -from unittest.mock import MagicMock +from unittest.mock import MagicMock, patch from langchain_core.documents import Document @@ -36,5 +37,23 @@ def test_duplicate_chunks_across_batches_are_submitted_once(self): self.assertEqual(collection.add_documents.call_count, 1) + def test_signature_memory_cap_is_respected_across_batches(self): + documents = [ + Document(page_content=content, metadata={"job_id": "job-1"}) + for content in ("alpha", "beta", "gamma", "alpha", "beta", "gamma") + ] + collection = MagicMock() + database = vector_db.__new__(vector_db) + database.collection = collection + + with patch.dict(os.environ, {"RAG_DEDUP_MAX_SIGNATURES": "2"}): + database.add_documents(documents, batch_size=3) + + submitted_batches = [ + len(call.args[0]) + for call in collection.add_documents.call_args_list + ] + self.assertEqual(submitted_batches, [3, 1]) + if __name__ == "__main__": unittest.main() From 1b8b60cabd1168ee8f556e842193b201954dc5ea Mon Sep 17 00:00:00 2001 From: Shantanav Mukherjee Date: Fri, 2 Oct 2026 00:43:56 +0530 Subject: [PATCH 8/8] Fix normalized whitespace in RAG fingerprinting --- python/app/rag/embeddings/reranker.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/app/rag/embeddings/reranker.py b/python/app/rag/embeddings/reranker.py index 773bc5e..5d8c668 100755 --- a/python/app/rag/embeddings/reranker.py +++ b/python/app/rag/embeddings/reranker.py @@ -5,7 +5,7 @@ from langchain_core.documents import Document -_WHITESPACE_RE = re.compile(r"\\s+") +_WHITESPACE_RE = re.compile(r"\s+") def normalize_text(text: str) -> str: