1use crate::{
2 api::{
3 ApiError, ApiMetadata, CodeRepositorySetAddResponse, CodeRepositorySetCreateResponse,
4 CodeRepositorySetQueryResponse, CodeRepositorySetRefreshResponse,
5 CodeRepositorySetStatusResponse, RequestContext,
6 },
7 code::{CodeIndexError, resolve_repository_snapshot},
8 domain::{
9 CodeQueryKind, CodeRepositorySelector, CodeRepositorySetAddMemberRequest,
10 CodeRepositorySetCreateRequest, CodeRepositorySetMemberStatus, CodeRepositorySetQueryHit,
11 CodeRepositorySetQueryRequest, CodeRepositorySetStatus, CodeRepositoryStatus,
12 CodeRetrievalHit, CodeRetrievalRequest, FreshnessPolicy,
13 },
14 storage::{
15 CodeRepositorySetMemberSeed, CodeRepositorySetRefreshTaskClaimRequest,
16 CodeRepositorySetRefreshTaskCompletion, CodeRepositorySetRefreshTaskFailure,
17 CodeRepositorySetRefreshTaskSeed, CodeRepositorySetSeed, StorageError,
18 },
19};
20use futures_util::{StreamExt, stream};
21use std::{path::PathBuf, sync::Arc};
22
23use super::{
24 RelayKnowledgeService,
25 code_repository_set_plan::{
26 dependency_symbol_plan_needs_hybrid_fallback, merge_dependency_symbol_fallback_hits,
27 repository_set_member_query_plan,
28 },
29 code_repository_set_query::{
30 OverlayEvidenceIndex, apply_bridge_support_bonus, dedupe_sort_truncate,
31 per_member_candidate_limit, prune_returned_overlay_evidence, repository_set_score,
32 },
33 code_service_support::apply_code_grep_fallback,
34};
35
36const REPOSITORY_SET_REFRESH_TASK_LEASE_MS: u64 = 10 * 60 * 1000;
37const REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS: u32 = 3;
38const REPOSITORY_SET_REFRESH_TASK_RETRY_BACKOFF_MS: u64 = 60_000;
39const REPOSITORY_SET_QUERY_MEMBER_CONCURRENCY: usize = 4;
40
41impl RelayKnowledgeService {
42 pub async fn create_code_repository_set(
44 &self,
45 request: CodeRepositorySetCreateRequest,
46 context: RequestContext,
47 ) -> Result<CodeRepositorySetCreateResponse, ApiError> {
48 let store = self.store().await.map_err(storage_api_error)?;
49 let repository_set = store
50 .create_code_repository_set(CodeRepositorySetSeed {
51 alias: request.alias.clone(),
52 description: request.description.clone(),
53 default_ref_policy_json: request.default_ref_policy_json.clone(),
54 now_ms: now_millis(),
55 })
56 .await
57 .map_err(storage_api_error)?;
58 let graph_version = store
59 .current_graph_version()
60 .await
61 .map_err(storage_api_error)?;
62
63 Ok(CodeRepositorySetCreateResponse {
64 metadata: ApiMetadata::graph_only(&context, graph_version),
65 request,
66 repository_set,
67 })
68 }
69
70 pub async fn add_code_repository_set_member(
72 &self,
73 request: CodeRepositorySetAddMemberRequest,
74 context: RequestContext,
75 ) -> Result<CodeRepositorySetAddResponse, ApiError> {
76 let store = self.store().await.map_err(storage_api_error)?;
77 let repository = store
78 .code_repository_status(request.repository_alias.clone())
79 .await
80 .map_err(storage_api_error)?
81 .ok_or_else(|| {
82 ApiError::invalid_argument(format!(
83 "code repository '{}' is not registered",
84 request.repository_alias
85 ))
86 })?;
87 let root_path = repository.root_path.clone();
88 let ref_selector = request.ref_selector.clone();
89 let (resolved_commit_sha, _tree_hash) =
90 run_blocking_code(move || resolve_repository_snapshot(root_path, &ref_selector))
91 .await?;
92 let path_filters = merged_filters(&repository.path_filters, &request.path_filters);
93 let language_filters =
94 merged_filters(&repository.language_filters, &request.language_filters);
95 let scope = store
96 .code_repository_scope_status(
97 request.repository_alias.clone(),
98 resolved_commit_sha.clone(),
99 path_filters.clone(),
100 language_filters.clone(),
101 )
102 .await
103 .map_err(storage_api_error)?
104 .ok_or_else(|| {
105 ApiError::invalid_argument(format!(
106 "code repository '{}' has no indexed scope for ref {} and requested filters",
107 request.repository_alias, request.ref_selector
108 ))
109 })?;
110 let source_scope = scope.last_indexed_scope_id.clone().ok_or_else(|| {
111 ApiError::invalid_argument(format!(
112 "code repository '{}' matching scope has no source scope",
113 request.repository_alias
114 ))
115 })?;
116 let member = store
117 .add_code_repository_set_member(CodeRepositorySetMemberSeed {
118 set_alias: request.set_alias.clone(),
119 repository_id: repository.repository_id,
120 repository_alias: request.repository_alias.clone(),
121 ref_selector: request.ref_selector.clone(),
122 resolved_commit_sha,
123 source_scope,
124 path_filters,
125 language_filters,
126 priority: request.priority,
127 })
128 .await
129 .map_err(storage_api_error)?;
130 let status = required_set_status(&store, &request.set_alias).await?;
131 let graph_version = store
132 .current_graph_version()
133 .await
134 .map_err(storage_api_error)?;
135
136 Ok(CodeRepositorySetAddResponse {
137 metadata: ApiMetadata::graph_only(&context, graph_version),
138 request,
139 member,
140 status,
141 })
142 }
143
144 pub async fn query_code_repository_set(
146 &self,
147 request: CodeRepositorySetQueryRequest,
148 context: RequestContext,
149 ) -> Result<CodeRepositorySetQueryResponse, ApiError> {
150 let store = self.store().await.map_err(storage_api_error)?;
151 let status = required_set_status(&store, &request.set_alias).await?;
152 let graph_version = store
153 .current_graph_version()
154 .await
155 .map_err(storage_api_error)?;
156 if request.freshness_policy == FreshnessPolicy::GraphOnly {
157 return Ok(CodeRepositorySetQueryResponse {
158 metadata: ApiMetadata::graph_only(&context, graph_version),
159 request,
160 status,
161 results: Vec::new(),
162 truncated: false,
163 degraded_reason: Some("graph_only freshness policy selected".to_owned()),
164 });
165 }
166 if let Some(error) = unfresh_set_error_for_wait_policy(&request, &status) {
167 return Err(error);
168 }
169 let edges = store
170 .code_repository_set_cross_edges(status.repository_set.set_id.clone())
171 .await
172 .map_err(storage_api_error)?;
173 let edge_index = OverlayEvidenceIndex::new(&edges);
174 let mut results = Vec::new();
175 let candidate_limit = per_member_candidate_limit(request.limit, status.members.len());
176 let highest_priority = status
177 .members
178 .iter()
179 .map(|member| member.member.priority)
180 .max()
181 .unwrap_or(0);
182 let member_outcomes = stream::iter(status.members.iter().cloned())
183 .map(|member_status| {
184 query_repository_set_member(
185 Arc::clone(&store),
186 request.clone(),
187 member_status,
188 highest_priority,
189 candidate_limit,
190 )
191 })
192 .buffer_unordered(REPOSITORY_SET_QUERY_MEMBER_CONCURRENCY)
193 .collect::<Vec<_>>()
194 .await;
195 let mut fallback_degraded_reasons = Vec::new();
196 for outcome in member_outcomes {
197 let outcome = outcome?;
198 fallback_degraded_reasons.push(outcome.degraded_reason);
199 for hit in outcome.hits {
200 let overlay_evidence = edge_index.evidence_for_hit(&hit);
201 let score = repository_set_score(&hit, &outcome.member_status, &overlay_evidence);
202 results.push(CodeRepositorySetQueryHit {
203 member: outcome.member_status.member.clone(),
204 hit,
205 overlay_evidence,
206 score,
207 });
208 }
209 }
210 apply_bridge_support_bonus(&mut results);
211 let truncated = dedupe_sort_truncate(&mut results, request.limit, &request.query);
212 prune_returned_overlay_evidence(&mut results);
213 let mut degraded_reasons = vec![
214 status.degraded_reason.clone(),
215 status
216 .overlay
217 .stale
218 .then(|| "repository set overlay is stale".to_owned()),
219 ];
220 degraded_reasons.extend(fallback_degraded_reasons);
221 let degraded_reason = join_degraded_reasons(degraded_reasons);
222
223 Ok(CodeRepositorySetQueryResponse {
224 metadata: ApiMetadata::graph_only(&context, graph_version),
225 request,
226 status,
227 results,
228 truncated,
229 degraded_reason,
230 })
231 }
232
233 pub async fn code_repository_set_status(
235 &self,
236 set_alias: String,
237 context: RequestContext,
238 ) -> Result<CodeRepositorySetStatusResponse, ApiError> {
239 let store = self.store().await.map_err(storage_api_error)?;
240 let status = required_set_status(&store, &set_alias).await?;
241 let graph_version = store
242 .current_graph_version()
243 .await
244 .map_err(storage_api_error)?;
245
246 Ok(CodeRepositorySetStatusResponse {
247 metadata: ApiMetadata::graph_only(&context, graph_version),
248 status,
249 })
250 }
251
252 pub async fn refresh_code_repository_set(
254 &self,
255 set_alias: String,
256 context: RequestContext,
257 ) -> Result<CodeRepositorySetRefreshResponse, ApiError> {
258 let store = self.store().await.map_err(storage_api_error)?;
259 let summary = store
260 .refresh_code_repository_set_overlay(set_alias.clone(), now_millis())
261 .await
262 .map_err(storage_api_error)?;
263 let status = required_set_status(&store, &set_alias).await?;
264 let graph_version = store
265 .current_graph_version()
266 .await
267 .map_err(storage_api_error)?;
268
269 Ok(CodeRepositorySetRefreshResponse {
270 metadata: ApiMetadata::graph_only(&context, graph_version),
271 status,
272 summary: Some(summary),
273 task: None,
274 })
275 }
276
277 pub async fn start_code_repository_set_refresh(
279 &self,
280 set_alias: String,
281 context: RequestContext,
282 ) -> Result<CodeRepositorySetRefreshResponse, ApiError> {
283 let store = self.store().await.map_err(storage_api_error)?;
284 let status = required_set_status(&store, &set_alias).await?;
285 let fingerprint = repository_set_refresh_fingerprint(&status);
286 let task = store
287 .queue_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskSeed {
288 set_id: status.repository_set.set_id.clone(),
289 set_alias: status.repository_set.alias.clone(),
290 input_fingerprint: fingerprint,
291 now_ms: now_millis(),
292 })
293 .await
294 .map_err(storage_api_error)?;
295 let graph_version = store
296 .current_graph_version()
297 .await
298 .map_err(storage_api_error)?;
299
300 Ok(CodeRepositorySetRefreshResponse {
301 metadata: ApiMetadata::graph_only(&context, graph_version),
302 status,
303 summary: None,
304 task: Some(task),
305 })
306 }
307
308 pub async fn run_code_repository_set_refresh_task_once(
310 &self,
311 task_id: Option<String>,
312 context: RequestContext,
313 ) -> Result<Option<crate::domain::CodeRepositorySetRefreshTaskRecord>, ApiError> {
314 let store = self.store().await.map_err(storage_api_error)?;
315 let lease_owner = format!("code-repository-set-refresh-worker-{}", std::process::id());
316 let Some(task) = store
317 .claim_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskClaimRequest {
318 task_id,
319 lease_owner: lease_owner.clone(),
320 lease_duration_ms: REPOSITORY_SET_REFRESH_TASK_LEASE_MS,
321 max_attempts: REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS,
322 now_ms: now_millis(),
323 })
324 .await
325 .map_err(storage_api_error)?
326 else {
327 return Ok(None);
328 };
329 let result = self
330 .refresh_code_repository_set(task.set_alias.clone(), context)
331 .await;
332 match result {
333 Ok(_) => store
334 .complete_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskCompletion {
335 task_id: task.task_id,
336 lease_owner,
337 attempt_count: task.attempt_count,
338 now_ms: now_millis(),
339 })
340 .await
341 .map(Some)
342 .map_err(storage_api_error),
343 Err(error) => {
344 let _ = store
345 .fail_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskFailure {
346 task_id: task.task_id,
347 lease_owner,
348 attempt_count: task.attempt_count,
349 error_kind: "repository_set_overlay".to_owned(),
350 error_message: error.message.clone(),
351 retry_backoff_ms: REPOSITORY_SET_REFRESH_TASK_RETRY_BACKOFF_MS,
352 max_attempts: REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS,
353 now_ms: now_millis(),
354 })
355 .await;
356 Err(error)
357 }
358 }
359 }
360
361 pub(crate) async fn code_repository_set_member_scopes(
362 &self,
363 set_alias: String,
364 ) -> Result<Option<Vec<(String, String)>>, ApiError> {
365 let store = self.store().await.map_err(storage_api_error)?;
366 store
367 .code_repository_set_status(set_alias)
368 .await
369 .map(|status| {
370 status.map(|status| {
371 status
372 .members
373 .into_iter()
374 .map(|member| (member.member.repository_alias, member.member.source_scope))
375 .collect()
376 })
377 })
378 .map_err(storage_api_error)
379 }
380}
381
382struct RepositorySetMemberQueryOutcome {
383 member_status: CodeRepositorySetMemberStatus,
384 hits: Vec<CodeRetrievalHit>,
385 degraded_reason: Option<String>,
386}
387
388async fn query_repository_set_member(
389 store: Arc<dyn crate::storage::KnowledgeStore>,
390 request: CodeRepositorySetQueryRequest,
391 member_status: CodeRepositorySetMemberStatus,
392 highest_priority: i32,
393 candidate_limit: usize,
394) -> Result<RepositorySetMemberQueryOutcome, ApiError> {
395 let member = &member_status.member;
396 let selector = CodeRepositorySelector::new(
397 member.repository_alias.clone(),
398 member.resolved_commit_sha.clone(),
399 request.path_filters.clone(),
400 request.language_filters.clone(),
401 )
402 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
403 let member_query_plan =
404 repository_set_member_query_plan(&request, &member_status, highest_priority);
405 let search_request = CodeRetrievalRequest::new(
406 member_query_plan.query,
407 selector.clone(),
408 member_query_plan.kind,
409 candidate_limit,
410 FreshnessPolicy::AllowStale,
411 )
412 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
413 let mut active_request = search_request.clone();
414 let mut hits = store
415 .search_code_scope(member.source_scope.clone(), search_request)
416 .await
417 .map_err(storage_api_error)?;
418 let dependency_symbol_plan_needs_fallback =
419 dependency_symbol_plan_needs_hybrid_fallback(&request, member_query_plan.kind, &hits);
420 let dependency_symbol_plan_satisfied = request.code_query_kind == CodeQueryKind::Hybrid
421 && member_query_plan.kind == CodeQueryKind::Symbol
422 && !dependency_symbol_plan_needs_fallback;
423 if dependency_symbol_plan_needs_fallback {
424 let symbol_plan_hits = hits;
425 let fallback_request = CodeRetrievalRequest::new(
426 request.query.clone(),
427 selector,
428 request.code_query_kind,
429 candidate_limit,
430 FreshnessPolicy::AllowStale,
431 )
432 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
433 active_request = fallback_request.clone();
434 let fallback_hits = store
435 .search_code_scope(member.source_scope.clone(), fallback_request)
436 .await
437 .map_err(storage_api_error)?;
438 hits = merge_dependency_symbol_fallback_hits(symbol_plan_hits, fallback_hits);
439 }
440 let degraded_reason = if repository_set_member_source_fallback_needed(
441 &request,
442 &active_request,
443 hits.len(),
444 dependency_symbol_plan_satisfied,
445 ) {
446 let base_status = required_member_repository(&store, &member.repository_id).await?;
447 let scoped_member_status =
448 code_status_for_repository_set_member(&base_status, &member_status);
449 apply_code_grep_fallback(
450 &store,
451 &base_status,
452 &scoped_member_status,
453 &active_request,
454 &mut hits,
455 )
456 .await?
457 } else {
458 None
459 };
460
461 Ok(RepositorySetMemberQueryOutcome {
462 member_status,
463 hits,
464 degraded_reason,
465 })
466}
467
468fn repository_set_member_source_fallback_needed(
469 set_request: &CodeRepositorySetQueryRequest,
470 active_request: &CodeRetrievalRequest,
471 hit_count: usize,
472 dependency_symbol_plan_satisfied: bool,
473) -> bool {
474 if dependency_symbol_plan_satisfied {
475 return false;
476 }
477
478 active_request.code_query_kind != CodeQueryKind::Hybrid || hit_count < set_request.limit.max(1)
479}
480
481pub(super) async fn required_set_status(
482 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
483 set_alias: &str,
484) -> Result<CodeRepositorySetStatus, ApiError> {
485 let mut status = store
486 .code_repository_set_status(set_alias.to_owned())
487 .await
488 .map_err(storage_api_error)?
489 .ok_or_else(|| {
490 ApiError::invalid_argument(format!(
491 "code repository set '{set_alias}' is not registered"
492 ))
493 })?;
494 refresh_moving_member_freshness(store, &mut status).await?;
495 refresh_repository_set_freshness(&mut status);
496
497 Ok(status)
498}
499
500fn join_degraded_reasons(reasons: impl IntoIterator<Item = Option<String>>) -> Option<String> {
501 let mut joined = Vec::new();
502 for reason in reasons.into_iter().flatten() {
503 if !joined.contains(&reason) {
504 joined.push(reason);
505 }
506 }
507
508 (!joined.is_empty()).then(|| joined.join("; "))
509}
510
511async fn refresh_moving_member_freshness(
512 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
513 status: &mut CodeRepositorySetStatus,
514) -> Result<(), ApiError> {
515 for index in 0..status.members.len() {
516 let member = status.members[index].member.clone();
517 let Some(reason) = moving_member_stale_reason(store, &member).await? else {
518 continue;
519 };
520 status.members[index].stale = true;
521 status.members[index].freshness_state = "stale".to_owned();
522 status.members[index].degraded_reason = Some(reason);
523 }
524
525 Ok(())
526}
527
528async fn moving_member_stale_reason(
529 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
530 member: &crate::domain::CodeRepositorySetMember,
531) -> Result<Option<String>, ApiError> {
532 if !member_ref_tracks_repository(&member.ref_selector, &member.resolved_commit_sha) {
533 return Ok(None);
534 }
535 let repository = store
536 .code_repository_status(member.repository_id.clone())
537 .await
538 .map_err(storage_api_error)?
539 .ok_or_else(|| {
540 ApiError::invalid_argument(format!(
541 "code repository '{}' is not registered",
542 member.repository_alias
543 ))
544 })?;
545 let root_path = PathBuf::from(repository.root_path);
546 let ref_selector = member.ref_selector.clone();
547 let resolved =
548 tokio::task::spawn_blocking(move || resolve_repository_snapshot(root_path, &ref_selector))
549 .await
550 .map_err(|error| ApiError::storage_unavailable(error.to_string()))?;
551
552 match resolved {
553 Ok((current_commit, _)) if current_commit == member.resolved_commit_sha => Ok(None),
554 Ok((current_commit, _)) => Ok(Some(format!(
555 "repository set member '{}' ref '{}' now resolves to {}, not stored snapshot {}",
556 member.repository_alias,
557 member.ref_selector,
558 current_commit,
559 member.resolved_commit_sha
560 ))),
561 Err(error) => Ok(Some(format!(
562 "repository set member '{}' ref '{}' could not be resolved: {error}",
563 member.repository_alias, member.ref_selector
564 ))),
565 }
566}
567
568fn member_ref_tracks_repository(ref_selector: &str, resolved_commit_sha: &str) -> bool {
569 let ref_selector = ref_selector.trim();
570 !(ref_selector == resolved_commit_sha
571 || (is_git_oid_prefix(ref_selector) && resolved_commit_sha.starts_with(ref_selector)))
572}
573
574fn is_git_oid_prefix(value: &str) -> bool {
575 (7..=64).contains(&value.len()) && value.bytes().all(|byte| byte.is_ascii_hexdigit())
576}
577
578fn refresh_repository_set_freshness(status: &mut CodeRepositorySetStatus) {
579 let member_stale = status.members.iter().any(|member| member.stale);
580 if member_stale && !status.overlay.stale {
581 status.overlay.stale = true;
582 status.overlay.state = "overlay_stale".to_owned();
583 }
584 status.freshness_state = if status.members.is_empty() {
585 "incomplete"
586 } else if member_stale {
587 "stale"
588 } else if status.overlay.stale {
589 "overlay_stale"
590 } else {
591 "fresh"
592 }
593 .to_owned();
594 status.degraded_reason = status
595 .members
596 .iter()
597 .find_map(|member| member.degraded_reason.clone())
598 .or_else(|| status.overlay.degraded_reason.clone());
599}
600
601fn unfresh_set_error_for_wait_policy(
602 request: &CodeRepositorySetQueryRequest,
603 status: &CodeRepositorySetStatus,
604) -> Option<ApiError> {
605 if request.freshness_policy != FreshnessPolicy::WaitUntilFresh {
606 return None;
607 }
608 if status.members.is_empty() {
609 return Some(ApiError::invalid_argument(format!(
610 "code repository set '{}' has no members",
611 status.repository_set.alias
612 )));
613 }
614 if let Some(member) = status.members.iter().find(|member| member.stale) {
615 return Some(ApiError::invalid_argument(format!(
616 "code repository set '{}' member '{}' scope '{}' is stale",
617 status.repository_set.alias, member.member.repository_alias, member.member.source_scope
618 )));
619 }
620 if status.overlay.stale {
621 return Some(ApiError::invalid_argument(format!(
622 "code repository set '{}' overlay is stale; run repo-set refresh before querying with wait_until_fresh",
623 status.repository_set.alias
624 )));
625 }
626
627 None
628}
629
630fn repository_set_refresh_fingerprint(status: &CodeRepositorySetStatus) -> String {
631 let mut parts = vec![status.repository_set.set_id.clone()];
632 parts.extend(status.members.iter().map(|member| {
633 format!(
634 "{}:{}:{}:{}:{}",
635 member.member.repository_id,
636 member.member.source_scope,
637 member.member.resolved_commit_sha,
638 member.tree_hash,
639 member.stale
640 )
641 }));
642 parts.join("|")
643}
644
645fn merged_filters(left: &[String], right: &[String]) -> Vec<String> {
646 let mut merged = Vec::new();
647 for value in left.iter().chain(right.iter()) {
648 if !merged.contains(value) {
649 merged.push(value.clone());
650 }
651 }
652
653 merged
654}
655
656async fn required_member_repository(
657 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
658 repository_id: &str,
659) -> Result<CodeRepositoryStatus, ApiError> {
660 store
661 .code_repository_status(repository_id.to_owned())
662 .await
663 .map_err(storage_api_error)?
664 .ok_or_else(|| {
665 ApiError::invalid_argument(format!(
666 "code repository set member repository '{repository_id}' is not registered"
667 ))
668 })
669}
670
671fn code_status_for_repository_set_member(
672 base_status: &CodeRepositoryStatus,
673 member_status: &CodeRepositorySetMemberStatus,
674) -> CodeRepositoryStatus {
675 let member = &member_status.member;
676 CodeRepositoryStatus {
677 repository_id: member.repository_id.clone(),
678 alias: member.repository_alias.clone(),
679 root_path: base_status.root_path.clone(),
680 path_filters: member.path_filters.clone(),
681 language_filters: member.language_filters.clone(),
682 last_indexed_scope_id: Some(member.source_scope.clone()),
683 last_indexed_commit: Some(member.resolved_commit_sha.clone()),
684 tree_hash: Some(member_status.tree_hash.clone()),
685 state: member_status.freshness_state.clone(),
686 indexed_file_count: member_status.indexed_file_count,
687 symbol_count: member_status.symbol_count,
688 reference_count: member_status.reference_count,
689 chunk_count: member_status.chunk_count,
690 stale: member_status.stale,
691 degraded_reason: member_status.degraded_reason.clone(),
692 }
693}
694
695fn now_millis() -> u64 {
696 std::time::SystemTime::now()
697 .duration_since(std::time::UNIX_EPOCH)
698 .map_or(0, |duration| {
699 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
700 })
701}
702
703async fn run_blocking_code<T, F>(operation: F) -> Result<T, ApiError>
704where
705 T: Send + 'static,
706 F: FnOnce() -> Result<T, CodeIndexError> + Send + 'static,
707{
708 tokio::task::spawn_blocking(operation)
709 .await
710 .map_err(|error| ApiError::storage_unavailable(error.to_string()))?
711 .map_err(code_api_error)
712}
713
714fn code_api_error(error: CodeIndexError) -> ApiError {
715 match error {
716 CodeIndexError::InvalidInput(message) => ApiError::invalid_argument(message),
717 CodeIndexError::Git { .. } | CodeIndexError::Io(_) | CodeIndexError::TreeSitter(_) => {
718 ApiError::storage_unavailable(error.to_string())
719 }
720 }
721}
722
723pub(super) fn storage_api_error(error: StorageError) -> ApiError {
724 match error {
725 StorageError::InvalidInput(message) => ApiError::invalid_argument(message),
726 other => ApiError::storage_unavailable(other.to_string()),
727 }
728}
729
730#[cfg(test)]
731mod tests {
732 use super::*;
733 use crate::{
734 api::ErrorKind,
735 domain::CodeRepositorySetMemberStatus,
736 domain::{CodeRepositorySet, CodeRepositorySetMember, CodeRepositorySetOverlayStatus},
737 storage::SqliteGraphStore,
738 };
739 use std::sync::Arc;
740
741 #[test]
742 fn helper_policy_reports_wait_until_fresh_blockers() {
743 let request = CodeRepositorySetQueryRequest::new(
744 "workspace",
745 "serve",
746 crate::domain::CodeQueryKind::Definition,
747 5,
748 FreshnessPolicy::WaitUntilFresh,
749 Vec::new(),
750 Vec::new(),
751 )
752 .expect("request should validate");
753 let empty = status_with_members(Vec::new(), overlay(true));
754 assert!(
755 unfresh_set_error_for_wait_policy(&request, &empty)
756 .expect("empty set should block")
757 .message
758 .contains("has no members")
759 );
760
761 let mut stale_member = member_status("app", "scope-app", 0);
762 stale_member.stale = true;
763 let stale_status = status_with_members(vec![stale_member], overlay(false));
764 assert!(
765 unfresh_set_error_for_wait_policy(&request, &stale_status)
766 .expect("stale member should block")
767 .message
768 .contains("member 'app'")
769 );
770
771 let overlay_status =
772 status_with_members(vec![member_status("app", "scope-app", 0)], overlay(true));
773 assert!(
774 unfresh_set_error_for_wait_policy(&request, &overlay_status)
775 .expect("stale overlay should block")
776 .message
777 .contains("overlay is stale")
778 );
779
780 let allow_stale = CodeRepositorySetQueryRequest::new(
781 "workspace",
782 "serve",
783 crate::domain::CodeQueryKind::Definition,
784 5,
785 FreshnessPolicy::AllowStale,
786 Vec::new(),
787 Vec::new(),
788 )
789 .expect("request should validate");
790 assert!(unfresh_set_error_for_wait_policy(&allow_stale, &overlay_status).is_none());
791 }
792
793 #[test]
794 fn helper_fingerprint_and_error_mapping_are_stable() {
795 let status = status_with_members(
796 vec![
797 member_status("app", "scope-app", 1),
798 member_status("svc", "scope-svc", 0),
799 ],
800 overlay(false),
801 );
802 let fingerprint = repository_set_refresh_fingerprint(&status);
803 assert!(fingerprint.contains("set-workspace"));
804 assert!(fingerprint.contains("repo-app:scope-app:commit-scope-app:tree-scope-app:false"));
805 assert_eq!(
806 merged_filters(&["src".to_owned()], &["src".to_owned(), "tests".to_owned()]),
807 ["src".to_owned(), "tests".to_owned()]
808 );
809 assert_eq!(
810 code_api_error(CodeIndexError::InvalidInput("bad ref".to_owned())).error_kind,
811 ErrorKind::InvalidArgument
812 );
813 assert_eq!(
814 code_api_error(CodeIndexError::Io(std::io::Error::other("disk"))).error_kind,
815 ErrorKind::StorageUnavailable
816 );
817 assert_eq!(
818 storage_api_error(StorageError::InvalidInput("bad storage".to_owned())).error_kind,
819 ErrorKind::InvalidArgument
820 );
821 }
822
823 #[test]
824 fn helper_source_fallback_policy_uses_final_set_limit_for_hybrid() {
825 let set_request = set_query_request(CodeQueryKind::Hybrid, 2);
826 let hybrid_member_request = member_retrieval_request(CodeQueryKind::Hybrid, 8);
827
828 assert!(!repository_set_member_source_fallback_needed(
829 &set_request,
830 &hybrid_member_request,
831 2,
832 false
833 ));
834 assert!(repository_set_member_source_fallback_needed(
835 &set_request,
836 &hybrid_member_request,
837 1,
838 false
839 ));
840 assert!(repository_set_member_source_fallback_needed(
841 &set_request,
842 &member_retrieval_request(CodeQueryKind::Imports, 8),
843 8,
844 false
845 ));
846 assert!(!repository_set_member_source_fallback_needed(
847 &set_request,
848 &member_retrieval_request(CodeQueryKind::Symbol, 8),
849 2,
850 true
851 ));
852 assert!(repository_set_member_source_fallback_needed(
853 &set_request,
854 &member_retrieval_request(CodeQueryKind::Symbol, 8),
855 2,
856 false
857 ));
858 }
859
860 #[tokio::test]
861 async fn helper_required_status_reports_missing_sets() {
862 let store: Arc<dyn crate::storage::KnowledgeStore> =
863 Arc::new(SqliteGraphStore::open_in_memory().expect("store should open"));
864 let error = required_set_status(&store, "missing")
865 .await
866 .expect_err("missing set should fail");
867
868 assert_eq!(error.error_kind, ErrorKind::InvalidArgument);
869 assert!(error.message.contains("is not registered"));
870 }
871
872 fn status_with_members(
873 members: Vec<CodeRepositorySetMemberStatus>,
874 overlay: CodeRepositorySetOverlayStatus,
875 ) -> CodeRepositorySetStatus {
876 CodeRepositorySetStatus {
877 repository_set: CodeRepositorySet {
878 set_id: "set-workspace".to_owned(),
879 alias: "workspace".to_owned(),
880 description: None,
881 default_ref_policy_json: "{\"default_ref\":\"HEAD\"}".to_owned(),
882 created_at_ms: 1,
883 updated_at_ms: 1,
884 },
885 members,
886 overlay,
887 freshness_state: "fresh".to_owned(),
888 degraded_reason: None,
889 }
890 }
891
892 fn set_query_request(kind: CodeQueryKind, limit: usize) -> CodeRepositorySetQueryRequest {
893 CodeRepositorySetQueryRequest::new(
894 "workspace",
895 "worker.New RegisterWorkflow",
896 kind,
897 limit,
898 FreshnessPolicy::AllowStale,
899 Vec::new(),
900 Vec::new(),
901 )
902 .expect("set query request should validate")
903 }
904
905 fn member_retrieval_request(kind: CodeQueryKind, limit: usize) -> CodeRetrievalRequest {
906 let selector =
907 CodeRepositorySelector::new("app", "commit", Vec::new(), Vec::new()).expect("selector");
908 CodeRetrievalRequest::new(
909 "worker.New RegisterWorkflow",
910 selector,
911 kind,
912 limit,
913 FreshnessPolicy::AllowStale,
914 )
915 .expect("member request should validate")
916 }
917
918 fn member_status(
919 repository_alias: &str,
920 source_scope: &str,
921 priority: i32,
922 ) -> CodeRepositorySetMemberStatus {
923 CodeRepositorySetMemberStatus {
924 member: CodeRepositorySetMember {
925 set_id: "set-workspace".to_owned(),
926 repository_id: format!("repo-{repository_alias}"),
927 repository_alias: repository_alias.to_owned(),
928 ref_selector: "HEAD".to_owned(),
929 resolved_commit_sha: format!("commit-{source_scope}"),
930 source_scope: source_scope.to_owned(),
931 path_filters: vec!["src".to_owned()],
932 language_filters: vec!["rust".to_owned()],
933 priority,
934 },
935 tree_hash: format!("tree-{source_scope}"),
936 freshness_state: "fresh".to_owned(),
937 stale: false,
938 indexed_file_count: 1,
939 symbol_count: 1,
940 reference_count: 0,
941 chunk_count: 1,
942 degraded_reason: None,
943 }
944 }
945
946 fn overlay(stale: bool) -> CodeRepositorySetOverlayStatus {
947 CodeRepositorySetOverlayStatus {
948 state: if stale { "overlay_stale" } else { "fresh" }.to_owned(),
949 stale,
950 edge_count: usize::from(!stale),
951 refreshed_at_ms: (!stale).then_some(10),
952 degraded_reason: None,
953 }
954 }
955}