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 mut 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 search_request.exclude_generated = request.exclude_generated;
463 let mut active_request = search_request.clone();
464 if let Some(reason) = fact_version_scope_mismatch_reason(&member_status) {
465 return Ok(RepositorySetMemberQueryOutcome {
466 member_status,
467 hits: Vec::new(),
468 active_request,
469 dependency_symbol_plan_satisfied: false,
470 source_fallback_allowed: false,
471 degraded_reason: Some(reason),
472 });
473 }
474 let mut hits = store
475 .search_code_scope(member.source_scope.clone(), search_request)
476 .await
477 .map_err(storage_api_error)?;
478 let dependency_symbol_plan_needs_fallback =
479 dependency_symbol_plan_needs_hybrid_fallback(&request, member_query_plan.kind, &hits);
480 let dependency_symbol_plan_satisfied = request.code_query_kind == CodeQueryKind::Hybrid
481 && member_query_plan.kind == CodeQueryKind::Symbol
482 && !dependency_symbol_plan_needs_fallback;
483 if dependency_symbol_plan_needs_fallback {
484 let symbol_plan_hits = hits;
485 let mut fallback_request = CodeRetrievalRequest::new(
486 request.query.clone(),
487 selector,
488 request.code_query_kind,
489 candidate_limit,
490 FreshnessPolicy::AllowStale,
491 )
492 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
493 fallback_request.exclude_generated = request.exclude_generated;
494 active_request = fallback_request.clone();
495 let fallback_hits = store
496 .search_code_scope(member.source_scope.clone(), fallback_request)
497 .await
498 .map_err(storage_api_error)?;
499 hits = merge_dependency_symbol_fallback_hits(symbol_plan_hits, fallback_hits);
500 }
501 Ok(RepositorySetMemberQueryOutcome {
502 member_status,
503 hits,
504 active_request,
505 dependency_symbol_plan_satisfied,
506 source_fallback_allowed: true,
507 degraded_reason: None,
508 })
509}
510
511fn repository_set_results_from_outcomes(
512 query: &str,
513 outcomes: &[RepositorySetMemberQueryOutcome],
514 edge_index: &OverlayEvidenceIndex<'_>,
515) -> Vec<CodeRepositorySetQueryHit> {
516 let mut results = Vec::new();
517 for outcome in outcomes {
518 for hit in &outcome.hits {
519 let overlay_evidence = edge_index.evidence_for_hit(hit);
520 let score = repository_set_score(query, hit, &outcome.member_status, &overlay_evidence);
521 results.push(CodeRepositorySetQueryHit {
522 member: outcome.member_status.member.clone(),
523 hit: hit.clone(),
524 overlay_evidence,
525 score,
526 });
527 }
528 }
529
530 results
531}
532
533fn repository_set_deferred_source_fallback_needed(
534 request: &CodeRepositorySetQueryRequest,
535 outcomes: &[RepositorySetMemberQueryOutcome],
536 initial_results: &[CodeRepositorySetQueryHit],
537) -> bool {
538 if outcomes.iter().any(|outcome| {
539 outcome.source_fallback_allowed
540 && outcome.active_request.code_query_kind != CodeQueryKind::Hybrid
541 && repository_set_member_source_fallback_needed(
542 request,
543 &outcome.active_request,
544 outcome.hits.len(),
545 outcome.dependency_symbol_plan_satisfied,
546 )
547 }) {
548 return true;
549 }
550 if outcomes.iter().any(|outcome| {
551 outcome.source_fallback_allowed
552 && outcome.hits.is_empty()
553 && repository_set_member_source_fallback_needed(
554 request,
555 &outcome.active_request,
556 outcome.hits.len(),
557 outcome.dependency_symbol_plan_satisfied,
558 )
559 }) {
560 return true;
561 }
562
563 let mut ranked = initial_results.to_vec();
564 dedupe_sort_truncate(&mut ranked, request.limit, &request.query);
565 ranked.len() < request.limit.max(1)
566}
567
568async fn apply_repository_set_deferred_source_fallbacks(
569 store: Arc<dyn crate::storage::KnowledgeStore>,
570 request: &CodeRepositorySetQueryRequest,
571 outcomes: &mut [RepositorySetMemberQueryOutcome],
572) -> Result<(), ApiError> {
573 let fallback_inputs = outcomes
574 .iter()
575 .enumerate()
576 .filter(|(_, outcome)| {
577 outcome.source_fallback_allowed
578 && repository_set_member_source_fallback_needed(
579 request,
580 &outcome.active_request,
581 outcome.hits.len(),
582 outcome.dependency_symbol_plan_satisfied,
583 )
584 })
585 .map(|(index, outcome)| RepositorySetMemberSourceFallbackInput {
586 index,
587 member_status: outcome.member_status.clone(),
588 active_request: outcome.active_request.clone(),
589 hits: outcome.hits.clone(),
590 })
591 .collect::<Vec<_>>();
592 let fallback_outputs = stream::iter(fallback_inputs)
593 .map(|input| {
594 let store = Arc::clone(&store);
595 async move { apply_repository_set_member_source_fallback(store, input).await }
596 })
597 .buffer_unordered(REPOSITORY_SET_QUERY_MEMBER_CONCURRENCY)
598 .collect::<Vec<_>>()
599 .await;
600 for output in fallback_outputs {
601 let output = output?;
602 outcomes[output.index].hits = output.hits;
603 outcomes[output.index].degraded_reason = output.degraded_reason;
604 }
605
606 Ok(())
607}
608
609async fn apply_repository_set_member_source_fallback(
610 store: Arc<dyn crate::storage::KnowledgeStore>,
611 input: RepositorySetMemberSourceFallbackInput,
612) -> Result<RepositorySetMemberSourceFallbackOutput, ApiError> {
613 let mut hits = input.hits;
614 let base_status =
615 required_member_repository(&store, &input.member_status.member.repository_id).await?;
616 let scoped_member_status =
617 code_status_for_repository_set_member(&base_status, &input.member_status);
618 let degraded_reason = apply_code_grep_fallback(
619 &store,
620 &base_status,
621 &scoped_member_status,
622 &input.active_request,
623 &mut hits,
624 )
625 .await?;
626
627 Ok(RepositorySetMemberSourceFallbackOutput {
628 index: input.index,
629 hits,
630 degraded_reason,
631 })
632}
633
634fn repository_set_member_source_fallback_needed(
635 set_request: &CodeRepositorySetQueryRequest,
636 active_request: &CodeRetrievalRequest,
637 hit_count: usize,
638 dependency_symbol_plan_satisfied: bool,
639) -> bool {
640 if dependency_symbol_plan_satisfied {
641 return false;
642 }
643
644 active_request.code_query_kind != CodeQueryKind::Hybrid || hit_count < set_request.limit.max(1)
645}
646
647pub(super) async fn required_set_status(
648 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
649 set_alias: &str,
650) -> Result<CodeRepositorySetStatus, ApiError> {
651 refreshed_required_set_status(store, set_alias)
652 .await
653 .map(|(status, _)| status)
654}
655
656async fn refreshed_required_set_status(
657 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
658 set_alias: &str,
659) -> Result<(CodeRepositorySetStatus, Vec<CodeRepositorySetMember>), ApiError> {
660 let mut status = store
661 .code_repository_set_status(set_alias.to_owned())
662 .await
663 .map_err(storage_api_error)?
664 .ok_or_else(|| {
665 ApiError::invalid_argument(format!(
666 "code repository set '{set_alias}' is not registered"
667 ))
668 })?;
669 let fact_version_replacements =
670 refresh_fact_version_member_freshness(store, &mut status).await?;
671 refresh_moving_member_freshness(store, &mut status).await?;
672 refresh_repository_set_freshness(&mut status);
673
674 Ok((status, fact_version_replacements))
675}
676
677async fn persist_fact_version_member_replacements(
678 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
679 set_alias: &str,
680 replacements: &[CodeRepositorySetMember],
681) -> Result<(), ApiError> {
682 for member in replacements {
683 store
684 .add_code_repository_set_member(CodeRepositorySetMemberSeed {
685 set_alias: set_alias.to_owned(),
686 repository_id: member.repository_id.clone(),
687 repository_alias: member.repository_alias.clone(),
688 ref_selector: member.ref_selector.clone(),
689 resolved_commit_sha: member.resolved_commit_sha.clone(),
690 source_scope: member.source_scope.clone(),
691 path_filters: member.path_filters.clone(),
692 language_filters: member.language_filters.clone(),
693 priority: member.priority,
694 })
695 .await
696 .map_err(storage_api_error)?;
697 }
698
699 Ok(())
700}
701
702fn join_degraded_reasons(reasons: impl IntoIterator<Item = Option<String>>) -> Option<String> {
703 let mut joined = Vec::new();
704 for reason in reasons.into_iter().flatten() {
705 if !joined.contains(&reason) {
706 joined.push(reason);
707 }
708 }
709
710 (!joined.is_empty()).then(|| joined.join("; "))
711}
712
713async fn refresh_moving_member_freshness(
714 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
715 status: &mut CodeRepositorySetStatus,
716) -> Result<(), ApiError> {
717 for index in 0..status.members.len() {
718 let member = status.members[index].member.clone();
719 let Some(reason) = moving_member_stale_reason(store, &member).await? else {
720 continue;
721 };
722 status.members[index].stale = true;
723 status.members[index].freshness_state = "stale".to_owned();
724 status.members[index].degraded_reason = Some(reason);
725 }
726
727 Ok(())
728}
729
730async fn moving_member_stale_reason(
731 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
732 member: &crate::domain::CodeRepositorySetMember,
733) -> Result<Option<String>, ApiError> {
734 if !member_ref_tracks_repository(&member.ref_selector, &member.resolved_commit_sha) {
735 return Ok(None);
736 }
737 let repository = store
738 .code_repository_status(member.repository_id.clone())
739 .await
740 .map_err(storage_api_error)?
741 .ok_or_else(|| {
742 ApiError::invalid_argument(format!(
743 "code repository '{}' is not registered",
744 member.repository_alias
745 ))
746 })?;
747 let ref_selector = member.ref_selector.clone();
748 let selector = CodeRepositorySelector {
749 repository: member.repository_alias.clone(),
750 ref_selector: ref_selector.clone(),
751 path_filters: member.path_filters.clone(),
752 language_filters: member.language_filters.clone(),
753 };
754 let resolved = resolve_code_ref_for_selector(&repository, &selector, ref_selector).await;
755
756 match resolved {
757 Ok(current_commit) if current_commit == member.resolved_commit_sha => Ok(None),
758 Ok(current_commit) => Ok(Some(format!(
759 "repository set member '{}' ref '{}' now resolves to {}, not stored snapshot {}",
760 member.repository_alias,
761 member.ref_selector,
762 current_commit,
763 member.resolved_commit_sha
764 ))),
765 Err(error) => Ok(Some(format!(
766 "repository set member '{}' ref '{}' could not be resolved: {error}",
767 member.repository_alias,
768 member.ref_selector,
769 error = error.message
770 ))),
771 }
772}
773
774fn member_ref_tracks_repository(ref_selector: &str, resolved_commit_sha: &str) -> bool {
775 let ref_selector = ref_selector.trim();
776 !(ref_selector == resolved_commit_sha
777 || (is_git_oid_prefix(ref_selector) && resolved_commit_sha.starts_with(ref_selector)))
778}
779
780fn is_git_oid_prefix(value: &str) -> bool {
781 (7..=64).contains(&value.len()) && value.bytes().all(|byte| byte.is_ascii_hexdigit())
782}
783
784fn refresh_repository_set_freshness(status: &mut CodeRepositorySetStatus) {
785 let member_stale = status.members.iter().any(|member| member.stale);
786 if member_stale && !status.overlay.stale {
787 status.overlay.stale = true;
788 status.overlay.state = "overlay_stale".to_owned();
789 }
790 status.freshness_state = if status.members.is_empty() {
791 "incomplete"
792 } else if member_stale {
793 "stale"
794 } else if status.overlay.stale {
795 "overlay_stale"
796 } else {
797 "fresh"
798 }
799 .to_owned();
800 status.degraded_reason = status
801 .members
802 .iter()
803 .find_map(|member| member.degraded_reason.clone())
804 .or_else(|| status.overlay.degraded_reason.clone());
805}
806
807fn unfresh_set_error_for_wait_policy(
808 request: &CodeRepositorySetQueryRequest,
809 status: &CodeRepositorySetStatus,
810) -> Option<ApiError> {
811 if request.freshness_policy != FreshnessPolicy::WaitUntilFresh {
812 return None;
813 }
814 if status.members.is_empty() {
815 return Some(ApiError::invalid_argument(format!(
816 "code repository set '{}' has no members",
817 status.repository_set.alias
818 )));
819 }
820 if let Some(member) = status.members.iter().find(|member| member.stale) {
821 return Some(ApiError::invalid_argument(format!(
822 "code repository set '{}' member '{}' scope '{}' is stale",
823 status.repository_set.alias, member.member.repository_alias, member.member.source_scope
824 )));
825 }
826 if status.overlay.stale {
827 return Some(ApiError::invalid_argument(format!(
828 "code repository set '{}' overlay is stale; run repo-set refresh before querying with wait_until_fresh",
829 status.repository_set.alias
830 )));
831 }
832
833 None
834}
835
836fn repository_set_refresh_fingerprint(status: &CodeRepositorySetStatus) -> String {
837 let mut parts = vec![status.repository_set.set_id.clone()];
838 parts.extend(status.members.iter().map(|member| {
839 format!(
840 "{}:{}:{}:{}:{}",
841 member.member.repository_id,
842 member.member.source_scope,
843 member.member.resolved_commit_sha,
844 member.tree_hash,
845 member.stale
846 )
847 }));
848 parts.join("|")
849}
850
851fn merged_filters(left: &[String], right: &[String]) -> Vec<String> {
852 let mut merged = Vec::new();
853 for value in left.iter().chain(right.iter()) {
854 if !merged.contains(value) {
855 merged.push(value.clone());
856 }
857 }
858
859 merged
860}
861
862async fn required_member_repository(
863 store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
864 repository_id: &str,
865) -> Result<CodeRepositoryStatus, ApiError> {
866 store
867 .code_repository_status(repository_id.to_owned())
868 .await
869 .map_err(storage_api_error)?
870 .ok_or_else(|| {
871 ApiError::invalid_argument(format!(
872 "code repository set member repository '{repository_id}' is not registered"
873 ))
874 })
875}
876
877fn code_status_for_repository_set_member(
878 base_status: &CodeRepositoryStatus,
879 member_status: &CodeRepositorySetMemberStatus,
880) -> CodeRepositoryStatus {
881 let member = &member_status.member;
882 CodeRepositoryStatus {
883 repository_id: member.repository_id.clone(),
884 alias: member.repository_alias.clone(),
885 root_path: base_status.root_path.clone(),
886 path_filters: member.path_filters.clone(),
887 language_filters: member.language_filters.clone(),
888 last_indexed_scope_id: Some(member.source_scope.clone()),
889 last_indexed_commit: Some(member.resolved_commit_sha.clone()),
890 tree_hash: Some(member_status.tree_hash.clone()),
891 state: member_status.freshness_state.clone(),
892 indexed_file_count: member_status.indexed_file_count,
893 symbol_count: member_status.symbol_count,
894 reference_count: member_status.reference_count,
895 chunk_count: member_status.chunk_count,
896 stale: member_status.stale,
897 degraded_reason: member_status.degraded_reason.clone(),
898 }
899}
900
901fn now_millis() -> u64 {
902 std::time::SystemTime::now()
903 .duration_since(std::time::UNIX_EPOCH)
904 .map_or(0, |duration| {
905 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
906 })
907}
908
909#[cfg(test)]
910fn code_api_error(error: CodeIndexError) -> ApiError {
911 match error {
912 CodeIndexError::InvalidInput(message) => ApiError::invalid_argument(message),
913 CodeIndexError::Git { .. } | CodeIndexError::Io(_) | CodeIndexError::TreeSitter(_) => {
914 ApiError::storage_unavailable(error.to_string())
915 }
916 }
917}
918
919pub(super) fn storage_api_error(error: StorageError) -> ApiError {
920 match error {
921 StorageError::InvalidInput(message) => ApiError::invalid_argument(message),
922 other => ApiError::storage_unavailable(other.to_string()),
923 }
924}
925
926#[cfg(test)]
927#[path = "service_tests.rs"]
928mod tests;