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