Skip to main content

relay_knowledge/application/
code_repository_set_service.rs

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