Coverage for src/local_deep_research/advanced_search_system/candidate_exploration/parallel_explorer.py: 99%

91 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-06 15:42 +0000

1""" 

2Parallel candidate explorer implementation. 

3 

4This explorer runs multiple search queries in parallel to quickly discover 

5a wide range of candidates. 

6""" 

7 

8import concurrent.futures 

9import time 

10from typing import List, Optional 

11 

12from loguru import logger 

13 

14from ...database.thread_local_session import thread_cleanup 

15from ...utilities.json_utils import get_llm_response_text 

16from ..candidates.base_candidate import Candidate 

17from ..constraints.base_constraint import Constraint 

18from .base_explorer import ( 

19 BaseCandidateExplorer, 

20 ExplorationResult, 

21 ExplorationStrategy, 

22) 

23 

24 

25class ParallelExplorer(BaseCandidateExplorer): 

26 """ 

27 Parallel candidate explorer that runs multiple searches concurrently. 

28 

29 This explorer: 

30 1. Generates multiple search queries from the initial query 

31 2. Runs searches in parallel for speed 

32 3. Collects and deduplicates candidates 

33 4. Focuses on breadth-first exploration 

34 """ 

35 

36 def __init__( 

37 self, 

38 *args, 

39 max_workers: int = 5, 

40 queries_per_round: int = 8, 

41 max_rounds: int = 3, 

42 **kwargs, 

43 ): 

44 """ 

45 Initialize parallel explorer. 

46 

47 Args: 

48 max_workers: Maximum number of parallel search threads 

49 queries_per_round: Number of queries to generate per round 

50 max_rounds: Maximum exploration rounds 

51 """ 

52 super().__init__(*args, **kwargs) 

53 self.max_workers = max_workers 

54 self.queries_per_round = queries_per_round 

55 self.max_rounds = max_rounds 

56 

57 def explore( 

58 self, 

59 initial_query: str, 

60 constraints: Optional[List[Constraint]] = None, 

61 entity_type: Optional[str] = None, 

62 ) -> ExplorationResult: 

63 """Explore candidates using parallel search strategy.""" 

64 start_time = time.time() 

65 logger.info(f"Starting parallel exploration for: {initial_query}") 

66 

67 all_candidates = [] 

68 exploration_paths = [] 

69 total_searched = 0 

70 

71 # Initial search 

72 current_queries = [initial_query] 

73 

74 with concurrent.futures.ThreadPoolExecutor( 

75 max_workers=self.max_workers 

76 ) as executor: 

77 for round_num in range(self.max_rounds): 

78 if not self._should_continue_exploration( 

79 start_time, len(all_candidates) 

80 ): 

81 break 

82 

83 logger.info( 

84 f"Exploration round {round_num + 1}: {len(current_queries)} queries" 

85 ) 

86 

87 # Submit all queries for parallel execution 

88 future_to_query = { 

89 executor.submit( 

90 thread_cleanup(self._execute_search), query 

91 ): query 

92 for query in current_queries 

93 } 

94 

95 round_candidates = [] 

96 

97 # Collect results as they complete 

98 for future in concurrent.futures.as_completed(future_to_query): 

99 query = future_to_query[future] 

100 total_searched += 1 

101 

102 try: 

103 results = future.result() 

104 candidates = self._extract_candidates_from_results( 

105 results, entity_type 

106 ) 

107 round_candidates.extend(candidates) 

108 exploration_paths.append( 

109 f"Round {round_num + 1}: {query} -> {len(candidates)} candidates" 

110 ) 

111 

112 except Exception: 

113 logger.exception(f"Error processing query '{query}'") 

114 

115 # Add new candidates 

116 all_candidates.extend(round_candidates) 

117 

118 # Generate queries for next round 

119 if round_num < self.max_rounds - 1: 

120 current_queries = self.generate_exploration_queries( 

121 initial_query, all_candidates, constraints 

122 )[: self.queries_per_round] 

123 

124 if not current_queries: 

125 logger.info("No more queries to explore") 

126 break 

127 

128 # Deduplicate and rank 

129 unique_candidates = self._deduplicate_candidates(all_candidates) 

130 ranked_candidates = self._rank_candidates_by_relevance( 

131 unique_candidates, initial_query 

132 ) 

133 

134 # Limit to max candidates 

135 final_candidates = ranked_candidates[: self.max_candidates] 

136 

137 elapsed_time = time.time() - start_time 

138 logger.info( 

139 f"Parallel exploration completed: {len(final_candidates)} unique candidates in {elapsed_time:.1f}s" 

140 ) 

141 

142 return ExplorationResult( 

143 candidates=final_candidates, 

144 total_searched=total_searched, 

145 unique_candidates=len(unique_candidates), 

146 exploration_paths=exploration_paths, 

147 metadata={ 

148 "strategy": "parallel", 

149 "rounds": min(round_num + 1, self.max_rounds), 

150 "max_workers": self.max_workers, 

151 "entity_type": entity_type, 

152 }, 

153 elapsed_time=elapsed_time, 

154 strategy_used=ExplorationStrategy.BREADTH_FIRST, 

155 ) 

156 

157 def generate_exploration_queries( 

158 self, 

159 base_query: str, 

160 found_candidates: List[Candidate], 

161 constraints: Optional[List[Constraint]] = None, 

162 ) -> List[str]: 

163 """Generate queries for parallel exploration.""" 

164 queries = [] 

165 

166 # Query variations based on base query 

167 base_variations = self._generate_query_variations(base_query) 

168 queries.extend(base_variations) 

169 

170 # Queries based on found candidates 

171 if found_candidates: 

172 candidate_queries = self._generate_candidate_based_queries( 

173 found_candidates, base_query 

174 ) 

175 queries.extend(candidate_queries) 

176 

177 # Constraint-based queries 

178 if constraints: 

179 constraint_queries = self._generate_constraint_queries( 

180 constraints, base_query 

181 ) 

182 queries.extend(constraint_queries) 

183 

184 # Remove already explored queries 

185 new_queries = [ 

186 q for q in queries if q.lower() not in self.explored_queries 

187 ] 

188 

189 return new_queries[: self.queries_per_round] 

190 

191 def _generate_query_variations(self, base_query: str) -> List[str]: 

192 """Generate variations of the base query.""" 

193 try: 

194 prompt = f""" 

195Generate 4 search query variations for: "{base_query}" 

196 

197Each variation should: 

1981. Use different keywords but same intent 

1992. Be specific and searchable 

2003. Focus on finding concrete examples or instances 

201 

202Format as numbered list: 

2031. [query] 

2042. [query] 

2053. [query] 

2064. [query] 

207""" 

208 

209 response = get_llm_response_text(self.model.invoke(prompt)).strip() 

210 

211 # Parse numbered list 

212 queries = [] 

213 for line in response.split("\n"): 

214 line = line.strip() 

215 if line and any(line.startswith(f"{i}.") for i in range(1, 10)): 

216 # Remove number prefix 

217 query = line.split(".", 1)[1].strip() 

218 if query: 218 ↛ 213line 218 didn't jump to line 213 because the condition on line 218 was always true

219 queries.append(query) 

220 

221 return queries[:4] 

222 

223 except Exception: 

224 logger.exception("Error generating query variations") 

225 return [] 

226 

227 def _generate_candidate_based_queries( 

228 self, candidates: List[Candidate], base_query: str 

229 ) -> List[str]: 

230 """Generate queries based on found candidates.""" 

231 queries = [] 

232 

233 # Sample a few candidates to avoid too many queries 

234 sample_candidates = candidates[:3] 

235 

236 for candidate in sample_candidates: 

237 # Query for similar entities 

238 queries.append(f'similar to "{candidate.name}"') 

239 queries.append(f'like "{candidate.name}" examples') 

240 

241 return queries 

242 

243 def _generate_constraint_queries( 

244 self, constraints: List[Constraint], base_query: str 

245 ) -> List[str]: 

246 """Generate queries focusing on specific constraints.""" 

247 queries = [] 

248 

249 # Sample constraints to avoid too many queries 

250 for constraint in constraints[:2]: 

251 queries.append(f"{constraint.value} examples") 

252 queries.append(f'"{constraint.value}" instances') 

253 

254 return queries