Repository navigation
Expand file tree
/
Copy pathflockparsecli.py
More file actions
2187 lines (1782 loc) · 83.2 KB
/
Copy pathflockparsecli.py
File metadata and controls
2187 lines (1782 loc) · 83.2 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
# CRITICAL: Enable unbuffered output FIRST to prevent CLI display freezing
# This ensures real-time progress messages are visible during async operations
import os
import sys
# Use development SOLLOL code instead of installed package
sys.path.insert(0, "/home/joker/SOLLOL/src")
# CRITICAL: Configure NetworkObserver BEFORE any SOLLOL imports
# NetworkObserver is a singleton initialized on first import, so this MUST come first
os.environ["SOLLOL_OBSERVER_SAMPLING"] = "false" # Disable sampling to get all events
os.environ["SOLLOL_REDIS_URL"] = "redis://localhost:6379" # Enable Redis pub/sub for dashboard
os.environ.setdefault("SOLLOL_ROUTING_LOG", "true") # Ensure routing decisions hit observability
# Force unbuffered output for stdout and stderr
# This prevents output buffering that can make the CLI appear frozen
os.environ["PYTHONUNBUFFERED"] = "1"
sys.stdout.reconfigure(line_buffering=True)
sys.stderr.reconfigure(line_buffering=True)
# CRITICAL: Configure logging BEFORE importing SOLLOL to prevent Dask logging deadlocks
from logging_config import setup_logging
logger = setup_logging()
# Pre-flight Redis connectivity check so observability issues are logged before SOLLOL loads
redis_url = os.getenv("SOLLOL_REDIS_URL", "redis://localhost:6379")
try:
import redis
from sollol.network_observer import logger as observer_logger
try:
redis.from_url(redis_url).ping()
msg = f"📡 Redis reachable at {redis_url}"
logger.info(msg)
observer_logger.info(msg)
except Exception as exc:
msg = f"⚠️ SOLLOL observability disabled: cannot reach Redis at {redis_url} ({exc})"
logger.warning(msg)
observer_logger.warning(msg)
except ImportError as exc:
logger.warning(f"⚠️ SOLLOL observability disabled: redis package not available ({exc})")
import ollama
# Set SOLLOL app name for logging context
os.environ["SOLLOL_APP_NAME"] = "FlockParser"
# Helper function to ensure prompts are visible
def visible_input(prompt):
"""Input with automatic stdout/stderr flushing to ensure prompt is visible."""
# Flush logger handlers first (now on stderr)
for handler in logger.handlers:
handler.flush()
# Flush both streams
sys.stdout.flush()
sys.stderr.flush()
# Write prompt to stderr (same stream as logger - known to be visible)
sys.stderr.write("\n" + prompt)
sys.stderr.flush()
# Read from stdin
return sys.stdin.readline().strip()
from pathlib import Path
from PyPDF2 import PdfReader
import docx
import subprocess
import tempfile
import json
import numpy as np
from datetime import datetime
import threading
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
import socket
import requests
import chromadb
from sollol.vram_monitor import VRAMMonitor, monitor_distributed_nodes
from gpu_controller import GPUController
from sollol.intelligent_gpu_router import IntelligentGPURouter
from sollol.adaptive_parallelism import AdaptiveParallelismStrategy
from sollol import OllamaPool # Direct SOLLOL integration
from sollol_compat import add_flockparser_methods # FlockParser compatibility layer
from parallel_embedder import embed_batch_parallel # Legacy parallel embedding
# 🚀 AVAILABLE COMMANDS:
COMMANDS = """
📖 open_pdf <file> → Process a single PDF file
📂 open_dir <dir> → Process all PDFs in a directory
💬 chat → Chat with processed PDFs
📊 list_docs → List all processed documents
🔍 check_deps → Check for required dependencies
🌐 discover_nodes → Auto-discover Ollama nodes on local network
➕ add_node <url> → Manually add an Ollama node (e.g., http://192.168.1.100:11434)
➖ remove_node <url> → Remove an Ollama node from the pool
📋 list_nodes → List all configured Ollama nodes
🔬 verify_models → Check which models are available on each node
⚖️ lb_stats → Show load balancer statistics
🎯 set_routing <strategy> → Set routing: adaptive, round_robin, least_loaded, lowest_latency
🖥️ vram_report → Show detailed VRAM usage report
🚀 force_gpu <model> → Force model to GPU on all capable nodes
🎯 gpu_status → Show intelligent GPU routing status
🧠 gpu_route <model> → Show routing decision for a model
🔧 gpu_optimize → Trigger intelligent GPU optimization
✅ gpu_check <model> → Check which nodes can fit a model
📚 gpu_models → List all known models and sizes
🗑️ unload_model <model> → Unload a specific model from memory
🧹 cleanup_models → Unload all non-priority models
🔀 parallelism_report → Show adaptive parallelism analysis
🧹 clear_cache → Clear embedding cache (keeps documents)
🗑️ clear_db → Clear ChromaDB vector store (removes all documents)
❌ exit → Quit the program
🌐 API Server: Automatically starts on port 8000 (http://localhost:8000)
Configure with: FLOCKPARSER_API=true/false, FLOCKPARSER_API_PORT=8000
"""
# 🔥 AI MODELS
EMBEDDING_MODEL = "mxbai-embed-large"
CHAT_MODEL = "qwen3:8b" # Fast and fits in available RAM (5.2 GB)
# 🚀 MODEL CACHING CONFIGURATION
# Keep models in VRAM for faster inference (prevents reloading)
EMBEDDING_KEEP_ALIVE = "1h" # Embedding model used frequently for chunking/search
CHAT_KEEP_ALIVE = "15m" # Chat model used less frequently
# 📊 RAG CONFIGURATION
# Retrieval settings for chat
RETRIEVAL_TOP_K = 10 # Number of chunks to retrieve (default: 10)
RETRIEVAL_MIN_SIMILARITY = 0.3 # Minimum similarity score (0.0-1.0)
CHUNKS_TO_SHOW = 10 # Number of source chunks to display (show all retrieved)
# Acceptable model variations (allows flexible matching)
ACCEPTABLE_EMBEDDING_MODELS = [
"mxbai-embed-large",
"mxbai-embed-large:latest",
"nomic-embed-text",
"nomic-embed-text:latest",
"all-minilm",
"all-minilm:latest",
"bge-large",
"bge-large:latest",
]
ACCEPTABLE_CHAT_MODELS = [
"llama3.1",
"llama3.1:8b",
"llama3.2",
"llama3.2:latest",
"llama3.2:3b",
"llama3",
"llama3:latest",
"llama3:8b",
"mistral",
"mistral:latest",
"mixtral",
"mixtral:latest",
"qwen",
"qwen:latest",
"qwen2.5",
"qwen3",
"qwen3:14b",
"qwen3:8b",
"qwen3:4b",
"gemma2:9b",
"phi3",
"deepseek-coder-v2",
"codellama:13b",
]
# 🌐 OLLAMA LOAD BALANCER CONFIGURATION
# SOLLOL auto-discovers all Ollama nodes on the network
# 📁 Directory setup
_SCRIPT_DIR = Path(__file__).parent
PROCESSED_DIR = _SCRIPT_DIR / "converted_files"
PROCESSED_DIR.mkdir(exist_ok=True)
# 📚 Knowledge Base (legacy JSON storage - kept for backwards compatibility)
KB_DIR = _SCRIPT_DIR / "knowledge_base"
KB_DIR.mkdir(exist_ok=True)
# 🗄️ ChromaDB Vector Store (production storage)
CHROMA_DB_DIR = _SCRIPT_DIR / "chroma_db_cli"
CHROMA_DB_DIR.mkdir(exist_ok=True)
# Initialize ChromaDB client and collection
chroma_client = chromadb.PersistentClient(path=str(CHROMA_DB_DIR))
chroma_collection = chroma_client.get_or_create_collection(
name="documents", metadata={"hnsw:space": "cosine"} # Use cosine similarity for better semantic search
)
# ============================================================
# SOLLOL Direct Integration (replaces ~1100 lines of custom code)
# Pure SOLLOL OllamaPool - no adapter layer
# Original implementation backed up to /tmp/old_loadbalancer_backup.py
# ============================================================
# Global reference (will be initialized in setup_load_balancer())
load_balancer = None
def setup_load_balancer():
"""Initialize SOLLOL pool with auto-discovery and dashboard.
Must be called from within if __name__ == '__main__': to avoid
multiprocessing issues with Dask worker spawning.
"""
global load_balancer
# CRITICAL: Initialize routing logger BEFORE OllamaPool (singleton pattern)
# This ensures OllamaPool gets a properly configured routing logger with Redis
from sollol.routing_logger import get_routing_logger, enable_console_routing_log
import redis as redis_lib
try:
routing_redis = redis_lib.from_url(
os.getenv("SOLLOL_REDIS_URL", "redis://localhost:6379"), decode_responses=True
)
routing_redis.ping() # Test connection
routing_logger = get_routing_logger(redis_client=routing_redis, console_output=False)
logger.info(
f"✅ Routing logger pre-initialized (enabled={routing_logger.enabled}, redis_available={routing_logger.redis_available})"
)
logger.info(f" 📡 Publishing routing decisions to: sollol:routing_events")
except Exception as e:
logger.warning(f"⚠️ Routing logger Redis configuration failed: {e}")
# Create logger anyway (will fall back to local-only mode)
routing_logger = get_routing_logger()
# Initialize SOLLOL pool with ROUND_ROBIN routing for balanced load distribution
# SOLLOL handles adaptive parallelism, intelligent routing, and distributed coordination
from sollol.routing_strategy import RoutingStrategy
load_balancer = OllamaPool(
nodes=None, # Auto-discover all Ollama nodes on network
routing_strategy=RoutingStrategy.ROUND_ROBIN, # Balanced round-robin distribution
exclude_localhost=True, # Use real IP instead of localhost
discover_all_nodes=True, # Scan full network for all nodes
app_name="FlockParser", # Identify as FlockParser in dashboard
enable_ray=False, # Skip Ray (single app, no cross-app coordination needed)
enable_dask=True, # Enable Dask for distributed batch processing (fixed stdin issue)
enable_gpu_redis=True, # Required for SOLLOL metrics publisher (latency, routing logs)
redis_host=os.getenv("SOLLOL_REDIS_HOST", "localhost"),
redis_port=int(os.getenv("SOLLOL_REDIS_PORT", "6379")),
register_with_dashboard=False, # Delay registration until after dashboard starts (see setup_dashboard)
)
# DEBUG: Check observer configuration for dashboard activity logging
from sollol.network_observer import get_observer
observer = get_observer()
logger.info(f"🔍 Observer Redis configured: {observer.redis_client is not None}")
logger.info(f"🔍 Observer total events so far: {observer.get_stats()['total_events']}")
# Force immediate flush of any pending dashboard events for real-time updates
if observer.redis_client:
observer._flush_dashboard_batch()
logger.info(
f"✅ Dashboard event flushing enabled (batch_size={observer.batch_size}, timeout={observer.batch_timeout}s)"
)
# Load primed performance stats if available
primed_stats_file = _SCRIPT_DIR / "sollol_primed_stats.json"
if primed_stats_file.exists():
try:
import json
with open(primed_stats_file, "r") as f:
primed_data = json.load(f)
if primed_data.get("priming_complete"):
# Merge primed stats into SOLLOL pool
if "node_performance" in primed_data:
load_balancer.stats["node_performance"] = primed_data["node_performance"]
logger.info("✅ Loaded primed performance stats (optimized distribution)")
# Show distribution preview
node_perf = primed_data["node_performance"]
if node_perf:
total_throughput = sum(p.get("batch_throughput", 0.5) for p in node_perf.values())
logger.info("📊 Configured workload distribution:")
for node_key, perf in node_perf.items():
throughput = perf.get("batch_throughput", 0.5)
pct = (throughput / total_throughput) * 100
logger.info(f" {node_key}: {pct:.1f}%")
except Exception as e:
logger.warning(f"⚠️ Could not load primed stats: {e}")
logger.info(" Run: python prime_sollol_performance.py")
# Add FlockParser compatibility methods
load_balancer = add_flockparser_methods(load_balancer, KB_DIR)
logger.info(
"Shim patch status: %s / %s / %s",
getattr(load_balancer, "_make_request", None).__qualname__,
getattr(load_balancer, "_make_streaming_request", None).__qualname__,
getattr(load_balancer, "_embed_batch_sequential", None).__qualname__,
)
return load_balancer
# Dashboard configuration (read from environment)
import os
_dashboard_enabled = os.getenv("FLOCKPARSER_DASHBOARD", "true").lower() in ("true", "1", "yes", "on")
_dashboard_port = int(os.getenv("FLOCKPARSER_DASHBOARD_PORT", "8080"))
def setup_dashboard():
"""Start SOLLOL unified dashboard after pool creation.
Must be called from within if __name__ == '__main__': to avoid
multiprocessing issues with Dask worker spawning.
"""
if not _dashboard_enabled:
logger.info("📊 Dashboard disabled (set FLOCKPARSER_DASHBOARD=true to enable)")
return
# Install Redis log publisher for distributed log streaming
from sollol.dashboard_service import install_redis_log_publisher
try:
install_redis_log_publisher()
logger.info("📡 Redis log publisher installed - logs streaming to dashboard")
except Exception as e:
logger.warning(f"⚠️ Redis log publisher failed: {e}")
# Check if dashboard already running, if not start dashboard_service
import requests
import subprocess
import threading
import time
dashboard_url = f"http://localhost:{_dashboard_port}"
dashboard_running = False
try:
response = requests.get(f"{dashboard_url}/api/applications", timeout=1)
if response.status_code == 200:
dashboard_running = True
logger.info(f"✅ Dashboard already running at {dashboard_url}")
except requests.exceptions.RequestException:
logger.info(f"🚀 Starting dashboard service at {dashboard_url}")
if not dashboard_running:
# Start dashboard_service as subprocess (not daemon thread)
dashboard_proc = subprocess.Popen(
[
"python3",
"-m",
"sollol.dashboard_service",
"--port",
str(_dashboard_port),
"--redis-url",
"redis://localhost:6379",
"--ray-dashboard-port",
"8265",
"--dask-dashboard-port",
"8787",
],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
# Wait for dashboard to start (Ray/Dask can take 10-15 seconds to initialize)
for attempt in range(30):
time.sleep(0.5)
try:
response = requests.get(f"{dashboard_url}/api/applications", timeout=1)
if response.status_code == 200:
logger.info(f"✅ Dashboard service started at {dashboard_url}")
dashboard_running = True
break
except requests.exceptions.RequestException:
continue
if not dashboard_running:
logger.error("❌ Dashboard failed to start after 15 seconds - check logs for errors")
# Dashboard registration - Register now that dashboard is running
# Use explicit DashboardClient registration (same as SynapticLlamas)
if dashboard_running and load_balancer:
try:
from sollol import DashboardClient
import socket
hostname = socket.gethostname()
dashboard_client = DashboardClient(
app_name=f"FlockParser ({hostname})",
router_type="OllamaPool",
version="1.0.0",
dashboard_url=dashboard_url,
metadata={
"nodes": len(load_balancer.nodes),
"routing_strategy": str(load_balancer.routing_strategy),
"enable_dask": True,
"enable_gpu_redis": True,
},
auto_register=True,
)
logger.info(f"✅ FlockParser registered with dashboard: {dashboard_client.app_id}")
except Exception as e:
logger.warning(f"⚠️ Dashboard registration failed: {e}")
logger.info(f"📊 Dashboard: {dashboard_url}")
logger.info(f" - Ray Dashboard: http://localhost:8265")
logger.info(f" - Dask Dashboard: http://localhost:8787")
# 💾 Index file for tracking processed documents
INDEX_FILE = KB_DIR / "document_index.json"
# 🔄 Cache for embeddings to avoid regenerating
EMBEDDING_CACHE_FILE = KB_DIR / "embedding_cache.json"
def load_embedding_cache():
"""Load the embedding cache from disk."""
if not EMBEDDING_CACHE_FILE.exists():
return {}
try:
with open(EMBEDDING_CACHE_FILE, "r") as f:
return json.load(f)
except (json.JSONDecodeError, FileNotFoundError):
return {}
def save_embedding_cache(cache):
"""Save the embedding cache to disk."""
with open(EMBEDDING_CACHE_FILE, "w") as f:
json.dump(cache, f)
def get_cached_embedding(text, use_load_balancer=True):
"""Get embedding from cache or generate new one."""
import hashlib
cache = load_embedding_cache()
# Create hash of text for cache key
text_hash = hashlib.md5(text.encode()).hexdigest()
if text_hash in cache:
return cache[text_hash]
# Generate new embedding using load balancer
if use_load_balancer:
embedding_result = load_balancer.embed(EMBEDDING_MODEL, text, keep_alive=EMBEDDING_KEEP_ALIVE, priority=7)
else:
embedding_result = ollama.embed(model=EMBEDDING_MODEL, input=text, keep_alive=EMBEDDING_KEEP_ALIVE)
embeddings = embedding_result.get("embeddings", [])
embedding = embeddings[0] if embeddings else []
# Cache it
cache[text_hash] = embedding
save_embedding_cache(cache)
return embedding
def load_document_index():
"""Load the document index or create it if it doesn't exist."""
if not INDEX_FILE.exists():
return {"documents": []}
try:
with open(INDEX_FILE, "r") as f:
return json.load(f)
except (json.JSONDecodeError, FileNotFoundError):
logger.error("⚠️ Error loading index file. Creating a new one.")
return {"documents": []}
def save_document_index(index_data):
"""Save the document index to disk."""
with open(INDEX_FILE, "w") as f:
json.dump(index_data, f, indent=4)
logger.info(f"✅ Document index updated with {len(index_data['documents'])} documents")
def register_document(pdf_path, txt_path, content, chunks=None):
"""Register a processed document in the knowledge base index."""
# Load existing index
index_data = load_document_index()
# Create document record
document_id = f"doc_{len(index_data['documents']) + 1}"
# Get PDF filename for better logging (especially in parallel mode)
from pathlib import Path
pdf_name = Path(pdf_path).stem if pdf_path else "unknown"
# Generate embeddings and chunks for search
chunks = chunks or chunk_text(content)
chunk_embeddings = []
# Batch process embeddings for better performance
logger.info(f"🔄 [{pdf_name}] Processing {len(chunks)} chunks in batches...")
import hashlib
cache = load_embedding_cache()
uncached_chunks = []
uncached_indices = []
# Check cache first
cached_count = 0
for i, chunk in enumerate(chunks):
text_hash = hashlib.md5(chunk.encode()).hexdigest()
if text_hash not in cache:
uncached_chunks.append(chunk)
uncached_indices.append(i)
else:
cached_count += 1
# Log cache status explicitly
if cached_count > 0:
logger.info(
f"📦 [{pdf_name}] Using {cached_count} cached embeddings, processing {len(uncached_chunks)} fresh chunks"
)
else:
logger.info(f"🆕 [{pdf_name}] No cached embeddings - processing all {len(uncached_chunks)} chunks fresh")
# Batch embed uncached chunks using SOLLOL's optimized embed_batch
if uncached_chunks:
logger.info(f"🚀 [{pdf_name}] Embedding {len(uncached_chunks)} new chunks...")
import time
start_time = time.time()
# Use SOLLOL's optimized embed_batch with adaptive parallelism
# This automatically:
# - Splits chunks across nodes (round-robin)
# - Uses connection reuse per node (_embed_batch_sequential)
# - Processes nodes in parallel
# - Provides 10-12x speedup over individual requests
logger.info(
f" 📞 Calling SOLLOL embed_batch with {len(uncached_chunks)} chunks, {len(load_balancer.nodes)} nodes"
)
# DEBUG: Check observer before embed_batch
from sollol.network_observer import get_observer
_obs = get_observer()
_before = _obs.get_stats()["total_events"]
logger.info(f"🔍 Observer events BEFORE embed_batch: {_before}")
all_results = load_balancer.embed_batch(
model=EMBEDDING_MODEL,
inputs=uncached_chunks,
priority=7,
use_adaptive=True, # Enable adaptive parallelism strategy
keep_alive=EMBEDDING_KEEP_ALIVE,
)
# DEBUG: Check observer after embed_batch
# CRITICAL: Observer processes events asynchronously in background thread!
# We need to wait for the event queue to be processed
import time
time.sleep(0.5) # Give background thread time to process events
_obs._flush_dashboard_batch() # Force flush any batched events
_after = _obs.get_stats()["total_events"]
logger.info(f"🔍 Observer events AFTER embed_batch (with 0.5s wait): {_after} (+{_after - _before})")
logger.info(f" 📥 SOLLOL embed_batch returned {len([r for r in all_results if r is not None])} results")
total_time = time.time() - start_time
successful = len([r for r in all_results if r is not None])
rate = successful / total_time if total_time > 0 else 0
logger.info(
f" ✅ Embedded {successful}/{len(uncached_chunks)} chunks in {total_time:.1f}s ({rate:.1f} chunks/sec)"
)
# Cache the embeddings
cached_count = 0
for chunk, result in zip(uncached_chunks, all_results):
if result:
text_hash = hashlib.md5(chunk.encode()).hexdigest()
embeddings = result.get("embeddings", [])
embedding = embeddings[0] if embeddings else []
cache[text_hash] = embedding
cached_count += 1
# Save cache once after all embeddings complete
save_embedding_cache(cache)
logger.info(f"✅ [{pdf_name}] Embedded and cached {cached_count}/{len(uncached_chunks)} chunks")
else:
logger.info(f"✅ [{pdf_name}] All chunks found in cache!")
# Now process all chunks
for i, chunk in enumerate(chunks):
try:
# Show progress every 50 chunks
if i % 50 == 0 and i > 0:
logger.info(f"🔄 Processed {i}/{len(chunks)} chunks...")
# Get embedding from cache
text_hash = hashlib.md5(chunk.encode()).hexdigest()
embedding = cache.get(text_hash, [])
# Store chunk with its embedding
chunk_file = KB_DIR / f"{document_id}_chunk_{i}.json"
chunk_data = {"text": chunk, "embedding": embedding}
with open(chunk_file, "w") as f:
json.dump(chunk_data, f)
# Remember the chunk reference
chunk_embeddings.append({"chunk_id": f"{document_id}_chunk_{i}", "file": str(chunk_file)})
except Exception as e:
logger.error(f"⚠️ Error embedding chunk {i}: {e}")
# Add document to index
doc_entry = {
"id": document_id,
"original": str(pdf_path),
"text_path": str(txt_path),
"processed_date": datetime.now().isoformat(),
"chunks": chunk_embeddings,
}
index_data["documents"].append(doc_entry)
save_document_index(index_data)
return document_id
def chunk_text(text, chunk_size=512, overlap=100):
"""
Split text into overlapping chunks with intelligent token-aware splitting.
Args:
chunk_size: Target chunk size in tokens (approximate via chars * 0.25)
overlap: Number of characters to overlap between chunks
"""
# Token limits for mxbai-embed-large: 512 tokens max
# Rough estimate: 1 token ≈ 4 characters
MAX_TOKENS = 480 # Leave buffer for model
MAX_CHARS = MAX_TOKENS * 4 # ~1920 chars
TARGET_CHARS = chunk_size * 4 # ~2048 chars for chunk_size=512
def split_large_text(text, max_size):
"""Recursively split text that's too large."""
if len(text) <= max_size:
return [text]
# Try splitting by sentences first
sentences = text.replace("! ", "!|").replace("? ", "?|").replace(". ", ".|").split("|")
chunks = []
current = []
current_len = 0
for sent in sentences:
sent = sent.strip()
if not sent:
continue
# If single sentence exceeds limit, split by words
if len(sent) > max_size:
words = sent.split()
# Calculate words per chunk (with safety margin)
words_per_chunk = int((max_size / len(sent)) * len(words) * 0.9)
words_per_chunk = max(50, words_per_chunk) # At least 50 words
for i in range(0, len(words), words_per_chunk):
word_chunk = " ".join(words[i : i + words_per_chunk])
if word_chunk:
chunks.append(word_chunk)
continue
# Add sentence to current chunk
if current_len + len(sent) > max_size and current:
chunks.append(" ".join(current))
current = [sent]
current_len = len(sent)
else:
current.append(sent)
current_len += len(sent)
if current:
chunks.append(" ".join(current))
return chunks
# Split into paragraphs first
paragraphs = [p.strip() for p in text.split("\n\n") if p.strip()]
chunks = []
current_chunk = []
current_length = 0
for para in paragraphs:
para_len = len(para)
# If paragraph is too large, split it first
if para_len > MAX_CHARS:
# Finalize current chunk if any
if current_chunk:
chunks.append("\n\n".join(current_chunk))
current_chunk = []
current_length = 0
# Split the large paragraph
para_chunks = split_large_text(para, MAX_CHARS)
chunks.extend(para_chunks)
continue
# Check if adding this paragraph exceeds target size
if current_length + para_len > TARGET_CHARS and current_chunk:
# Finalize current chunk
chunks.append("\n\n".join(current_chunk))
# Start new chunk with overlap (keep last paragraph if small enough)
if overlap > 0 and current_chunk and len(current_chunk[-1]) < overlap:
current_chunk = [current_chunk[-1], para]
current_length = len(current_chunk[-1]) + para_len
else:
current_chunk = [para]
current_length = para_len
else:
current_chunk.append(para)
current_length += para_len
# Add final chunk
if current_chunk:
final_chunk = "\n\n".join(current_chunk)
# Safety check
if len(final_chunk) > MAX_CHARS:
chunks.extend(split_large_text(final_chunk, MAX_CHARS))
else:
chunks.append(final_chunk)
# Final validation: ensure no chunk exceeds MAX_CHARS
validated_chunks = []
for chunk in chunks:
if len(chunk) > MAX_CHARS:
validated_chunks.extend(split_large_text(chunk, MAX_CHARS))
else:
validated_chunks.append(chunk)
return validated_chunks
def list_documents():
"""List all processed documents in the knowledge base."""
index_data = load_document_index()
if not index_data["documents"]:
logger.info("📚 No documents have been processed yet.")
return
logger.info(f"\n📚 Knowledge Base: {len(index_data['documents'])} documents")
logger.info("-" * 60)
for i, doc in enumerate(index_data["documents"]):
logger.info(f"{i+1}. {Path(doc['original']).name}")
logger.info(f" ID: {doc['id']} | Processed: {doc['processed_date'][:10]}")
logger.info(f" Chunks: {len(doc['chunks'])}")
logger.info("-" * 60)
def get_similar_chunks(query, top_k=None, min_similarity=None):
"""Find text chunks similar to the query using vector similarity with adaptive top-k."""
# Use configured defaults if not specified
if min_similarity is None:
min_similarity = RETRIEVAL_MIN_SIMILARITY
try:
# Get embedding for the query from cache
query_embedding = get_cached_embedding(query)
if not query_embedding:
logger.error("⚠️ Failed to generate query embedding")
return []
# Load document index
index_data = load_document_index()
# Check if we have documents
if not index_data["documents"]:
logger.info("📚 No documents in knowledge base yet")
return []
# Adaptive top-k based on total chunks in database
if top_k is None:
total_chunks = sum(len(doc["chunks"]) for doc in index_data["documents"])
# Scale top_k based on database size
if total_chunks < 50:
adaptive_k = min(total_chunks, 5) # Very small DB, use fewer
elif total_chunks < 200:
adaptive_k = 10 # Small-medium DB, use default
elif total_chunks < 1000:
adaptive_k = 20 # Medium DB, retrieve more context
else:
adaptive_k = 30 # Large DB, need more chunks for good coverage
top_k = adaptive_k
logger.info(f" 📊 Adaptive top-k: {top_k} (from {total_chunks} total chunks)")
else:
logger.info(f" 📊 Using fixed top-k: {top_k}")
# Collect all chunks with their embeddings
chunks_with_similarity = []
for doc in index_data["documents"]:
for chunk_ref in doc["chunks"]:
try:
# Load chunk data
chunk_file = Path(chunk_ref["file"])
if chunk_file.exists():
with open(chunk_file, "r") as f:
chunk_data = json.load(f)
# Calculate cosine similarity
chunk_embedding = chunk_data.get("embedding", [])
if chunk_embedding:
similarity = cosine_similarity(query_embedding, chunk_embedding)
if similarity >= min_similarity:
chunks_with_similarity.append(
{
"doc_id": doc["id"],
"doc_name": Path(doc["original"]).name,
"text": chunk_data["text"],
"similarity": similarity,
}
)
except Exception as e:
logger.error(f"⚠️ Error processing chunk {chunk_ref['chunk_id']}: {e}")
# Sort by similarity (highest first) and get top k
chunks_with_similarity.sort(key=lambda x: x["similarity"], reverse=True)
# Return top k results
results = chunks_with_similarity[:top_k]
# Print retrieval stats
logger.info(f" Found {len(results)} relevant chunks (similarity >= {min_similarity:.2f})")
return results
except Exception as e:
logger.error(f"⚠️ Error searching knowledge base: {e}")
return []
def sanitize_for_xml(text):
"""Remove null bytes and control characters that break XML/DOCX."""
import re
# Remove NULL bytes
text = text.replace("\x00", "")
# Remove other control characters except newline, carriage return, and tab
text = re.sub(r"[\x00-\x08\x0B-\x0C\x0E-\x1F\x7F-\x9F]", "", text)
return text
def cosine_similarity(vec1, vec2):
"""Calculate cosine similarity between two vectors."""
if not vec1 or not vec2:
return 0
vec1 = np.array(vec1)
vec2 = np.array(vec2)
dot_product = np.dot(vec1, vec2)
norm_a = np.linalg.norm(vec1)
norm_b = np.linalg.norm(vec2)
if norm_a == 0 or norm_b == 0:
return 0
return dot_product / (norm_a * norm_b)
def embed_text(text):
"""Embeds text using Ollama without storing vector data in files."""
try:
# Using 'input' instead of 'prompt'
_ = ollama.embed(model=EMBEDDING_MODEL, input=text)
return text # Return the original text for saving to files
except Exception as e:
logger.error(f"❌ Embedding error: {e}")
return None
def clean_extracted_text(text):
"""Clean extracted text by normalizing Unicode and fixing common LaTeX/PDF extraction issues."""
import re
import unicodedata
if not text:
return text
# Step 1: Normalize Unicode (convert composed chars to decomposed and back)
text = unicodedata.normalize("NFKC", text)
# Step 2: Fix common Unicode escape sequences that appear as literal text
# Replace \uXXXX patterns with actual Unicode characters
def replace_unicode_escapes(match):
try:
code = match.group(1)
return chr(int(code, 16))
except:
return match.group(0)
text = re.sub(r"\\u([0-9a-fA-F]{4})", replace_unicode_escapes, text)
text = re.sub(r"\\x([0-9a-fA-F]{2})", replace_unicode_escapes, text)
# Step 3: Clean up common LaTeX remnants that get corrupted
# Replace common Greek letter codes with their actual Unicode
greek_map = {
r"\\alpha": "α",
r"\\beta": "β",
r"\\gamma": "γ",
r"\\delta": "δ",
r"\\epsilon": "ε",
r"\\zeta": "ζ",
r"\\eta": "η",
r"\\theta": "θ",
r"\\iota": "ι",
r"\\kappa": "κ",
r"\\lambda": "λ",
r"\\mu": "μ",
r"\\nu": "ν",
r"\\xi": "ξ",
r"\\pi": "π",
r"\\rho": "ρ",
r"\\sigma": "σ",
r"\\tau": "τ",
r"\\upsilon": "υ",
r"\\phi": "φ",
r"\\chi": "χ",
r"\\psi": "ψ",
r"\\omega": "ω",
# Capital letters
r"\\Gamma": "Γ",
r"\\Delta": "Δ",
r"\\Theta": "Θ",
r"\\Lambda": "Λ",
r"\\Xi": "Ξ",
r"\\Pi": "Π",
r"\\Sigma": "Σ",
r"\\Phi": "Φ",
r"\\Psi": "Ψ",
r"\\Omega": "Ω",
}
for latex, unicode_char in greek_map.items():
text = text.replace(latex, unicode_char)
# Step 4: Fix spacing issues - add space after periods if missing
text = re.sub(r"\.([A-Z])", r". \1", text)
# Step 5: Remove excessive whitespace
text = re.sub(r"[ \t]+", " ", text) # Multiple spaces to single space
text = re.sub(r"\n{3,}", "\n\n", text) # Multiple newlines to double newline
return text.strip()
def extract_text_from_pdf(pdf_path):
"""Extracts text from a PDF file using multiple methods for better reliability."""
pdf_path_str = str(pdf_path)
extracted_text = ""
# Method 1: Try PyMuPDF (fitz) first - better word spacing preservation
try:
import fitz # PyMuPDF
logger.info("🔍 Attempting extraction with PyMuPDF (better word spacing)...")
doc = fitz.open(pdf_path_str)
pymupdf_text = ""
for page_num, page in enumerate(doc):
# extract_text() with "text" mode preserves word spacing better
page_text = page.get_text("text")
if page_text:
# Clean the text immediately after extraction
page_text = clean_extracted_text(page_text)
pymupdf_text += f"{page_text}\n\n"
else:
logger.warning(f"⚠️ PyMuPDF: No text extracted from page {page_num + 1}")
doc.close()
if pymupdf_text.strip():
logger.info(f"✅ PyMuPDF successfully extracted {len(pymupdf_text)} characters")
extracted_text = pymupdf_text
else:
logger.warning("⚠️ PyMuPDF extraction yielded no text, trying alternative method...")
except ImportError: