Skip to main content

relay_knowledge/application/code_repository/repository_set/
service.rs

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    /// Creates or updates a thin repository set.
48    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    /// Adds one already-indexed repository snapshot to a repository set.
76    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    /// Queries every member scope and merges ranked candidates without changing single-repo search.
154    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    /// Returns repository-set freshness and member diagnostics.
244    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    /// Rebuilds cross-repository import/module overlay edges.
263    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    /// Queues a repository-set overlay refresh task.
288    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    /// Runs one queued repository-set overlay refresh task under a lease.
319    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;