@@ -173,7 +173,13 @@ def _resolve_paths(
173173 self ,
174174 paths : Optional [Union [str , Path , List [str ], List [Path ]]],
175175 ) -> List [str ]:
176- """Resolve and normalise paths: arg > self.paths > cwd.
176+ """Resolve and normalise paths with layered fallback.
177+
178+ Priority (highest → lowest):
179+ 1. Explicit ``paths`` argument (``search(..., paths=xxx)``)
180+ 2. Instance default ``self.paths`` (constructor ``paths=``)
181+ 3. ``SIRCHMUNK_SEARCH_PATHS`` environment variable (comma-separated)
182+ 4. Current working directory
177183
178184 Always returns ``List[str]`` so callers need no further coercion.
179185 """
@@ -183,6 +189,14 @@ def _resolve_paths(
183189 return [str (p ) for p in paths ]
184190 if self .paths is not None :
185191 return list (self .paths )
192+ env_paths = os .getenv ("SIRCHMUNK_SEARCH_PATHS" , "" )
193+ if env_paths :
194+ parsed = [p .strip () for p in env_paths .split ("," ) if p .strip ()]
195+ if parsed :
196+ _loguru_logger .info (
197+ f"[paths] Using SIRCHMUNK_SEARCH_PATHS: { parsed } "
198+ )
199+ return parsed
186200 cwd = str (Path .cwd ())
187201 _loguru_logger .info (
188202 f"[paths] No paths provided; using current working directory: { cwd } "
@@ -285,81 +299,56 @@ def _load_historical_knowledge(self):
285299 async def _try_reuse_cluster (self , query : str ) -> Optional [KnowledgeCluster ]:
286300 """Try to reuse existing knowledge cluster based on semantic similarity.
287301
302+ The method waits (non-blocking) for the embedding model to become
303+ ready so that reuse works reliably even on the first search call
304+ within a process.
305+
288306 Returns:
289307 KnowledgeCluster if a suitable cached cluster is found, None otherwise.
290308 """
291309 if not self .embedding_client :
292310 return None
293311
294- # Skip cluster reuse while the embedding model is still loading in
295- # its background thread; kick off loading so it's ready next time.
296- if not self .embedding_client .is_ready ():
297- self .embedding_client .start_loading ()
298- return None
299-
300312 try :
313+ # Wait for the model (non-blocking via executor) instead of
314+ # returning None immediately — this ensures reuse works even
315+ # on the very first search call.
316+ if not self .embedding_client .is_ready ():
317+ self .embedding_client .start_loading ()
318+ try :
319+ await self .embedding_client ._ensure_model_async (timeout = 60 )
320+ except Exception :
321+ await self ._logger .debug (
322+ "Embedding model not ready yet, skipping cluster reuse"
323+ )
324+ return None
325+
301326 await self ._logger .info ("Searching for similar knowledge clusters..." )
302-
303- # Compute query embedding
327+
304328 query_embedding = (await self .embedding_client .embed ([query ]))[0 ]
305-
306- # Search for similar clusters
329+
307330 similar_clusters = await self .knowledge_storage .search_similar_clusters (
308331 query_embedding = query_embedding ,
309332 top_k = self .cluster_sim_top_k ,
310333 similarity_threshold = self .cluster_sim_threshold ,
311334 )
312-
335+
313336 if not similar_clusters :
314337 await self ._logger .info ("No similar clusters found, performing new search..." )
315338 return None
316-
317- # Found similar cluster - process reuse
339+
318340 best_match = similar_clusters [0 ]
319341 await self ._logger .success (
320342 f"♻️ Found similar cluster: { best_match ['name' ]} "
321343 f"(similarity: { best_match ['similarity' ]:.3f} )"
322344 )
323-
324- # Retrieve full cluster object
345+
325346 existing_cluster = await self .knowledge_storage .get (best_match ["id" ])
326-
327347 if not existing_cluster :
328348 await self ._logger .warning ("Failed to retrieve cluster, falling back to new search" )
329349 return None
330-
331- # Add current query to queries list with FIFO strategy
332- self ._add_query_to_cluster (existing_cluster , query )
333-
334- # Update hotness and timestamp for reused cluster
335- existing_cluster .hotness = min (1.0 , (existing_cluster .hotness or 0.5 ) + 0.1 )
336- existing_cluster .last_modified = datetime .now ()
337-
338- # Recompute embedding with new query (before update to avoid double save)
339- if self .embedding_client and self .embedding_client .is_ready ():
340- try :
341- from sirchmunk .utils .embedding_util import compute_text_hash
342350
343- combined_text = self .knowledge_storage .combine_cluster_fields (
344- existing_cluster .queries
345- )
346- text_hash = compute_text_hash (combined_text )
347- embedding_vector = (await self .embedding_client .embed ([combined_text ]))[0 ]
348-
349- await self .knowledge_storage .store_embedding (
350- cluster_id = existing_cluster .id ,
351- embedding_vector = embedding_vector ,
352- embedding_model = self .embedding_client .model_id ,
353- embedding_text_hash = text_hash ,
354- )
355- await self ._logger .debug (f"Updated embedding for cluster { existing_cluster .id } " )
356- except Exception as emb_error :
357- await self ._logger .warning (f"Failed to update embedding: { emb_error } " )
358-
359- # Single update call - saves cluster data and embedding together
360- await self .knowledge_storage .update (existing_cluster )
361-
362- # Validate cluster has usable content
351+ # Validate cluster has usable content BEFORE mutating it
363352 content = existing_cluster .content
364353 if isinstance (content , list ):
365354 content = "\n " .join (content )
@@ -369,9 +358,41 @@ async def _try_reuse_cluster(self, query: str) -> Optional[KnowledgeCluster]:
369358 )
370359 return None
371360
361+ # Mutate only after validation passes
362+ self ._add_query_to_cluster (existing_cluster , query )
363+ existing_cluster .hotness = min (1.0 , (existing_cluster .hotness or 0.5 ) + 0.1 )
364+ existing_cluster .last_modified = datetime .now ()
365+
366+ # Recompute embedding with updated queries list
367+ try :
368+ from sirchmunk .utils .embedding_util import compute_text_hash
369+
370+ combined_text = self .knowledge_storage .combine_cluster_fields (
371+ existing_cluster .queries
372+ )
373+ text_hash = compute_text_hash (combined_text )
374+ embedding_vector = (await self .embedding_client .embed ([combined_text ]))[0 ]
375+
376+ await self .knowledge_storage .store_embedding (
377+ cluster_id = existing_cluster .id ,
378+ embedding_vector = embedding_vector ,
379+ embedding_model = self .embedding_client .model_id ,
380+ embedding_text_hash = text_hash ,
381+ )
382+ except Exception as emb_error :
383+ await self ._logger .warning (f"Failed to update embedding: { emb_error } " )
384+
385+ await self .knowledge_storage .update (existing_cluster )
386+
387+ # Flush to parquet so the updated cluster is visible to future searches
388+ try :
389+ self .knowledge_storage .force_sync ()
390+ except Exception as sync_err :
391+ await self ._logger .warning (f"Parquet force_sync failed: { sync_err } " )
392+
372393 await self ._logger .success ("Reused existing knowledge cluster" )
373394 return existing_cluster
374-
395+
375396 except Exception as e :
376397 await self ._logger .warning (
377398 f"Failed to search similar clusters: { e } . Falling back to full search."
@@ -406,12 +427,17 @@ async def _save_cluster_with_embedding(self, cluster: KnowledgeCluster) -> None:
406427 Args:
407428 cluster: KnowledgeCluster to save
408429 """
409- # Save knowledge cluster to persistent storage
430+ # Save knowledge cluster to persistent storage.
431+ # insert() returns False (without raising) when the cluster already
432+ # exists, so we explicitly fall back to update() in that case.
410433 try :
411- await self .knowledge_storage .insert (cluster )
412- await self ._logger .info (f"Saved knowledge cluster { cluster .id } to cache" )
434+ inserted = await self .knowledge_storage .insert (cluster )
435+ if inserted :
436+ await self ._logger .info (f"Saved knowledge cluster { cluster .id } to cache" )
437+ else :
438+ await self .knowledge_storage .update (cluster )
439+ await self ._logger .info (f"Updated knowledge cluster { cluster .id } in cache" )
413440 except Exception as e :
414- # If cluster exists, update it instead
415441 try :
416442 await self .knowledge_storage .update (cluster )
417443 await self ._logger .info (f"Updated knowledge cluster { cluster .id } in cache" )
@@ -978,6 +1004,8 @@ async def _search_deep(
9781004
9791005 # ==============================================================
9801006 # Phase 0: Cluster reuse (instant short-circuit)
1007+ # When reuse_knowledge=True and a similar cluster is found, we
1008+ # return here — Phase 5 (Persistence) is not executed for that path.
9811009 # ==============================================================
9821010 reused = await self ._try_reuse_cluster (query )
9831011 if reused is not None :
@@ -1116,14 +1144,19 @@ async def _search_deep(
11161144 context .add_llm_tokens (total_tok , usage = usage )
11171145
11181146 # ==============================================================
1119- # Phase 5: Persistence
1147+ # Phase 5: Persistence (only when no cluster was reused in Phase 0)
1148+ # When Phase 0 reuses a cluster we return early, so this block
1149+ # runs only for newly built clusters from Phase 1–4.
11201150 # ==============================================================
11211151 phase5_tasks = []
11221152 if cluster :
11231153 self ._add_query_to_cluster (cluster , query )
11241154 phase5_tasks .append (self ._save_cluster_with_embedding (cluster ))
11251155 phase5_tasks .append (self ._save_spec_context (paths , context , scan_result = scan_result ))
1126- await asyncio .gather (* phase5_tasks , return_exceptions = True )
1156+ results = await asyncio .gather (* phase5_tasks , return_exceptions = True )
1157+ for r in results :
1158+ if isinstance (r , Exception ):
1159+ _loguru_logger .warning (f"[Phase 5] Persistence task failed: { r } " )
11271160
11281161 await self ._logger .success (f"[search] Complete: { context .summary ()} " )
11291162 return answer , cluster , context
@@ -1219,6 +1252,7 @@ async def _search_fast(
12191252
12201253 # ==============================================================
12211254 # Step 0: Cluster reuse — instant short-circuit (no LLM cost)
1255+ # When reuse succeeds we return here; no persistence step runs.
12221256 # ==============================================================
12231257 reused = await self ._try_reuse_cluster (query )
12241258 if reused is not None :
@@ -1336,7 +1370,7 @@ async def _search_fast(
13361370 query , answer , file_path , evidence , keywords_used ,
13371371 )
13381372
1339- # Persist the FAST cluster so it can be reused by future queries
1373+ # Persist the new cluster (only reached when Step 0 did not reuse)
13401374 self ._add_query_to_cluster (cluster , query )
13411375 try :
13421376 await self ._save_cluster_with_embedding (cluster )
0 commit comments