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