Skip to main content

relay_knowledge/application/code_repository/
repository.rs

1use std::{collections::BTreeMap, path::PathBuf};
2
3use crate::{
4    api::{
5        ApiError, ApiMetadata, CodeRepositoryFeatureFlagsResponse, CodeRepositoryImpactResponse,
6        CodeRepositoryIndexResponse, CodeRepositoryIndexStartResponse, CodeRepositoryQueryResponse,
7        CodeRepositoryRegisterRequest, CodeRepositoryRegisterResponse,
8        CodeRepositoryRemoveResponse, CodeRepositoryReportResponse,
9        CodeRepositoryScopePreviewResponse, CodeRepositoryStatusResponse, RequestContext,
10    },
11    code::{
12        REGISTRATION_LANGUAGE_FILTER_ERROR, build_index_snapshot_with_workspace_detection,
13        changed_paths_for_diff_with_filters, changed_paths_for_filesystem_diff,
14        deleted_symbol_names_for_diff, partition_changed_paths_for_selector,
15        prepare_full_index_plan_with_workspace_detection, preview_repository_scope,
16        register_repository,
17    },
18    domain::{
19        CodeFeatureFlagRequest, CodeImpactRequest, CodeIndexMode, CodeIndexRequest,
20        CodeIndexResourceBudget, CodeRepositorySelector, CodeRepositoryStatus,
21        CodeRetrievalRequest, FreshnessPolicy,
22    },
23    storage::CodeImpactChanges,
24};
25
26#[path = "repository_staleness.rs"]
27mod repository_staleness;
28pub(super) use repository_staleness::annotate_query_result_staleness;
29
30use crate::application::service::RelayKnowledgeService;
31
32use super::fast_index::fresh_full_index_response;
33use super::freshness::{
34    CodeFeatureFlagFreshnessContext, CodeQueryFreshnessContext,
35    code_feature_flag_freshness_diagnostics, code_query_freshness_diagnostics,
36};
37use super::support::{
38    CODE_INDEX_TASK_LEASE_MS, CODE_INDEX_TASK_MAX_ATTEMPTS, CODE_INDEX_TASK_RETRY_BACKOFF_MS,
39    CodeIndexTaskLeaseContext, RETAIN_RECENT_CODE_SCOPES, active_index_matches_request,
40    apply_code_grep_fallback, code_index_worker_lease_owner, code_status_checkpoint,
41    feature_flag_request_at_indexed_ref, index_start_from_completed, indexed_source_scope,
42    latest_compatible_code_scope_status, missing_indexed_source_scope_error, now_millis,
43    previous_index_state_for_index, recover_code_index_task_leases,
44    recover_orphaned_code_index_task_leases, refresh_code_index_task_lease,
45    registration_from_status, required_code_repository, resolve_code_ref_for_selector,
46    resolved_code_scope_status, retrieval_request_at_indexed_ref, run_blocking_code,
47    storage_api_error,
48};
49
50impl RelayKnowledgeService {
51    /// Registers a Git repository as a code source.
52    pub async fn register_code_repository(
53        &self,
54        request: CodeRepositoryRegisterRequest,
55        context: RequestContext,
56    ) -> Result<CodeRepositoryRegisterResponse, ApiError> {
57        if !request.language_filters.is_empty() {
58            return Err(ApiError::invalid_argument(
59                REGISTRATION_LANGUAGE_FILTER_ERROR,
60            ));
61        }
62        let registration = run_blocking_code(move || {
63            register_repository(
64                request.root_path,
65                request.alias,
66                request.path_filters,
67                request.language_filters,
68            )
69        })
70        .await?;
71        let store = self.store().await.map_err(storage_api_error)?;
72        let status = store
73            .upsert_code_repository(registration.clone())
74            .await
75            .map_err(storage_api_error)?;
76        let _ = self.refresh_watched_code_repository(&status).await;
77        let graph_version = store
78            .current_graph_version()
79            .await
80            .map_err(storage_api_error)?;
81
82        Ok(CodeRepositoryRegisterResponse {
83            metadata: ApiMetadata::graph_only(&context, graph_version),
84            registration,
85            status,
86        })
87    }
88
89    /// Removes a registered code repository and its derived index state.
90    pub async fn remove_code_repository(
91        &self,
92        repository: String,
93        context: RequestContext,
94    ) -> Result<CodeRepositoryRemoveResponse, ApiError> {
95        let store = self.store().await.map_err(storage_api_error)?;
96        let now_ms = now_millis();
97        recover_code_index_task_leases(&store, now_ms).await?;
98        let removed_status = required_code_repository(&store, &repository).await?;
99        let summary = store
100            .remove_code_repository(removed_status.repository_id.clone(), now_ms)
101            .await
102            .map_err(storage_api_error)?
103            .ok_or_else(|| {
104                ApiError::storage_unavailable("removed code repository disappeared before delete")
105            })?;
106        let _ = self
107            .remove_watched_code_repository(&removed_status.alias, &removed_status.repository_id)
108            .await;
109        let graph_version = store
110            .current_graph_version()
111            .await
112            .map_err(storage_api_error)?;
113
114        Ok(CodeRepositoryRemoveResponse {
115            metadata: ApiMetadata::graph_only(&context, graph_version),
116            removed_status,
117            summary,
118        })
119    }
120
121    /// Builds or updates the tree-sitter code index for a registered repository.
122    pub async fn index_code_repository(
123        &self,
124        request: CodeIndexRequest,
125        context: RequestContext,
126    ) -> Result<CodeRepositoryIndexResponse, ApiError> {
127        self.index_code_repository_inner(request, context, None)
128            .await
129    }
130
131    async fn index_code_repository_inner(
132        &self,
133        request: CodeIndexRequest,
134        context: RequestContext,
135        task_lease: Option<CodeIndexTaskLeaseContext>,
136    ) -> Result<CodeRepositoryIndexResponse, ApiError> {
137        let store = self.store().await.map_err(storage_api_error)?;
138        let status = required_code_repository(&store, &request.repository.repository).await?;
139        if let Some(response) =
140            fresh_full_index_response(&store, &status, &request, &context).await?
141        {
142            return Ok(response);
143        }
144        let registration = registration_from_status(&status);
145        let selector = request.repository.clone();
146        let summary = if request.mode == CodeIndexMode::Full {
147            self.apply_full_code_index(
148                &store,
149                registration,
150                selector,
151                request.workspace_detection.clone(),
152                CodeIndexResourceBudget::default(),
153                task_lease,
154            )
155            .await?
156        } else {
157            let previous = previous_index_state_for_index(&store, &status, &request).await?;
158            let mode = request.mode;
159            let workspace_detection = request.workspace_detection.clone();
160            let snapshot = run_blocking_code(move || {
161                build_index_snapshot_with_workspace_detection(
162                    &registration,
163                    &selector,
164                    mode,
165                    previous.fingerprints,
166                    previous.base_resolved_commit_sha,
167                    &workspace_detection,
168                )
169            })
170            .await?;
171            store
172                .apply_code_index_snapshot(snapshot)
173                .await
174                .map_err(storage_api_error)?
175        };
176        let status = store
177            .code_repository_status(summary.repository_id.clone())
178            .await
179            .map_err(storage_api_error)?
180            .ok_or_else(|| ApiError::storage_unavailable("code repository status is missing"))?;
181        let software_projection = store
182            .refresh_software_global_projection(summary.source_scope.clone())
183            .await
184            .map_err(storage_api_error)?;
185        let graph_version = store
186            .current_graph_version()
187            .await
188            .map_err(storage_api_error)?;
189        let degraded_reason = status
190            .degraded_reason
191            .clone()
192            .or(software_projection.status.last_error.clone());
193        let status = CodeRepositoryStatus {
194            degraded_reason,
195            ..status
196        };
197        let _ = self.refresh_watched_code_repository(&status).await;
198
199        Ok(CodeRepositoryIndexResponse {
200            metadata: ApiMetadata::graph_only(&context, graph_version),
201            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
202                &status,
203                &request.repository,
204                request.repository.ref_selector.clone(),
205            ),
206            summary,
207            status,
208        })
209    }
210
211    async fn apply_full_code_index(
212        &self,
213        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
214        registration: crate::domain::CodeRepositoryRegistration,
215        selector: CodeRepositorySelector,
216        workspace_detection: crate::domain::CodeWorkspaceDetectionConfig,
217        resource_budget: CodeIndexResourceBudget,
218        task_lease: Option<CodeIndexTaskLeaseContext>,
219    ) -> Result<crate::domain::CodeIndexSummary, ApiError> {
220        let mut plan = run_blocking_code(move || {
221            prepare_full_index_plan_with_workspace_detection(
222                registration,
223                selector,
224                resource_budget,
225                &workspace_detection,
226            )
227        })
228        .await?;
229        let session = plan.session();
230        store
231            .begin_code_index_session(session.clone())
232            .await
233            .map_err(storage_api_error)?;
234        loop {
235            refresh_code_index_task_lease(store, task_lease.as_ref()).await?;
236            let (next_plan, batch) = run_blocking_code(move || plan.parse_next_batch()).await?;
237            plan = next_plan;
238            let Some(batch) = batch else {
239                break;
240            };
241            store
242                .apply_code_index_batch(batch)
243                .await
244                .map_err(storage_api_error)?;
245            refresh_code_index_task_lease(store, task_lease.as_ref()).await?;
246        }
247
248        refresh_code_index_task_lease(store, task_lease.as_ref()).await?;
249        let summary = store
250            .finalize_code_index_session(session)
251            .await
252            .map_err(storage_api_error)?;
253        refresh_code_index_task_lease(store, task_lease.as_ref()).await?;
254
255        Ok(summary)
256    }
257
258    /// Starts a repository index request, queueing cold full indexes for background execution.
259    pub async fn start_code_repository_index(
260        &self,
261        request: CodeIndexRequest,
262        context: RequestContext,
263    ) -> Result<CodeRepositoryIndexStartResponse, ApiError> {
264        let store = self.store().await.map_err(storage_api_error)?;
265        let status = required_code_repository(&store, &request.repository.repository).await?;
266        if let Some(response) =
267            fresh_full_index_response(&store, &status, &request, &context).await?
268        {
269            return Ok(index_start_from_completed(response, None));
270        }
271        if request.mode != CodeIndexMode::Full {
272            let response = self.index_code_repository(request, context).await?;
273            return Ok(index_start_from_completed(response, None));
274        }
275        recover_code_index_task_leases(&store, now_millis()).await?;
276        let payload_json = serde_json::to_string(&request)
277            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
278        if let Some(active_task) = self
279            .active_full_index_task_for_request(&store, &status, &request, &payload_json)
280            .await?
281        {
282            return self
283                .index_start_response_from_task(
284                    &store,
285                    status,
286                    active_task,
287                    request.repository.ref_selector,
288                    &context,
289                )
290                .await;
291        }
292
293        let registration = registration_from_status(&status);
294        let selector = request.repository.clone();
295        let workspace_detection = request.workspace_detection.clone();
296        let workspace_detection_json = serde_json::to_string(&workspace_detection)
297            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
298        let resource_budget = CodeIndexResourceBudget::default();
299        let plan = run_blocking_code(move || {
300            prepare_full_index_plan_with_workspace_detection(
301                registration,
302                selector,
303                resource_budget,
304                &workspace_detection,
305            )
306        })
307        .await?;
308        let session = plan.session();
309        let input_fingerprint = format!(
310            "full:{}:{}:{}:{}",
311            session.repository_id,
312            session.tree_hash,
313            session.source_scope,
314            workspace_detection_json
315        );
316        if let Some(active_task) = store
317            .active_code_index_task(session.repository_id.clone())
318            .await
319            .map_err(storage_api_error)?
320            && active_task.state.is_unfinished()
321            && active_task.input_fingerprint == input_fingerprint
322        {
323            return self
324                .index_start_response_from_task(
325                    &store,
326                    status,
327                    active_task,
328                    request.repository.ref_selector,
329                    &context,
330                )
331                .await;
332        }
333        let task = store
334            .queue_code_index_task(crate::storage::CodeIndexTaskSeed {
335                repository_id: session.repository_id.clone(),
336                alias: status.alias.clone(),
337                ref_selector: request.repository.ref_selector.clone(),
338                resolved_commit_sha: session.resolved_commit_sha.clone(),
339                tree_hash: session.tree_hash.clone(),
340                source_scope: session.source_scope.clone(),
341                path_filters: session.path_filters.clone(),
342                language_filters: session.language_filters.clone(),
343                mode: request.mode.clone(),
344                input_fingerprint,
345                resource_budget: session.resource_budget,
346                payload_json,
347                now_ms: now_millis(),
348            })
349            .await
350            .map_err(storage_api_error)?;
351        self.index_start_response_from_task(
352            &store,
353            status,
354            task,
355            request.repository.ref_selector,
356            &context,
357        )
358        .await
359    }
360
361    async fn active_full_index_task_for_request(
362        &self,
363        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
364        status: &CodeRepositoryStatus,
365        request: &CodeIndexRequest,
366        payload_json: &str,
367    ) -> Result<Option<crate::domain::CodeIndexTaskRecord>, ApiError> {
368        let Some(active_task) = store
369            .active_code_index_task(status.repository_id.clone())
370            .await
371            .map_err(storage_api_error)?
372        else {
373            return Ok(None);
374        };
375        if !active_task.state.is_unfinished()
376            || active_task.mode != CodeIndexMode::Full
377            || active_task.payload_json != payload_json
378        {
379            return Ok(None);
380        }
381        let resolved = resolve_code_ref_for_selector(
382            status,
383            &request.repository,
384            request.repository.ref_selector.clone(),
385        )
386        .await?;
387        if resolved == active_task.resolved_commit_sha {
388            Ok(Some(active_task))
389        } else {
390            Ok(None)
391        }
392    }
393
394    async fn index_start_response_from_task(
395        &self,
396        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
397        fallback_status: CodeRepositoryStatus,
398        task: crate::domain::CodeIndexTaskRecord,
399        requested_ref: String,
400        context: &RequestContext,
401    ) -> Result<CodeRepositoryIndexStartResponse, ApiError> {
402        let checkpoint = store
403            .code_index_checkpoint(task.source_scope.clone())
404            .await
405            .map_err(storage_api_error)?;
406        let graph_version = store
407            .current_graph_version()
408            .await
409            .map_err(storage_api_error)?;
410        let status = store
411            .code_repository_status(task.repository_id.clone())
412            .await
413            .map_err(storage_api_error)?
414            .unwrap_or(fallback_status);
415
416        Ok(CodeRepositoryIndexStartResponse {
417            metadata: ApiMetadata::graph_only(context, graph_version),
418            scope: crate::api::CodeRepositoryScopeMetadata::from_index_task(&task, requested_ref),
419            summary: None,
420            status,
421            task: Some(task),
422            checkpoint,
423        })
424    }
425
426    /// Runs one queued code index task under a lease.
427    pub async fn run_code_index_task_once(
428        &self,
429        task_id: Option<String>,
430        context: RequestContext,
431    ) -> Result<Option<crate::domain::CodeIndexTaskRecord>, ApiError> {
432        let store = self.store().await.map_err(storage_api_error)?;
433        let lease_owner = code_index_worker_lease_owner();
434        let Some(task) = store
435            .claim_code_index_task(crate::storage::CodeIndexTaskClaimRequest {
436                task_id,
437                lease_owner: lease_owner.clone(),
438                lease_duration_ms: CODE_INDEX_TASK_LEASE_MS,
439                max_attempts: CODE_INDEX_TASK_MAX_ATTEMPTS,
440                now_ms: now_millis(),
441            })
442            .await
443            .map_err(storage_api_error)?
444        else {
445            return Ok(None);
446        };
447        let mut request = match serde_json::from_str::<CodeIndexRequest>(&task.payload_json) {
448            Ok(request) => request,
449            Err(error) => {
450                let message = format!(
451                    "code index task '{}' payload is invalid: {error}",
452                    task.task_id
453                );
454                let _ = store
455                    .fail_code_index_task(crate::storage::CodeIndexTaskFailure {
456                        task_id: task.task_id,
457                        lease_owner,
458                        attempt_count: task.attempt_count,
459                        error_kind: "task_payload".to_owned(),
460                        error_message: message.clone(),
461                        retry_backoff_ms: CODE_INDEX_TASK_RETRY_BACKOFF_MS,
462                        max_attempts: CODE_INDEX_TASK_MAX_ATTEMPTS,
463                        now_ms: now_millis(),
464                    })
465                    .await;
466                return Err(ApiError::invalid_argument(message));
467            }
468        };
469        if task.mode != CodeIndexMode::WorktreeOverlay {
470            request.repository.ref_selector = task.resolved_commit_sha.clone();
471        }
472        let lease_context = CodeIndexTaskLeaseContext {
473            task_id: task.task_id.clone(),
474            lease_owner: lease_owner.clone(),
475            attempt_count: task.attempt_count,
476            lease_duration_ms: CODE_INDEX_TASK_LEASE_MS,
477        };
478        let result = self
479            .index_code_repository_inner(request, context, Some(lease_context.clone()))
480            .await;
481        match result {
482            Ok(response) => {
483                refresh_code_index_task_lease(&store, Some(&lease_context)).await?;
484                let completed = store
485                    .complete_code_index_task(crate::storage::CodeIndexTaskCompletion {
486                        task_id: task.task_id.clone(),
487                        lease_owner,
488                        attempt_count: task.attempt_count,
489                        now_ms: now_millis(),
490                    })
491                    .await
492                    .map_err(storage_api_error)?;
493                let _ = store
494                    .prune_code_repository_scopes(crate::storage::CodeScopeRetentionRequest {
495                        repository_id: response.summary.repository_id,
496                        active_scope: response.summary.source_scope,
497                        retain_recent_successful_scopes: RETAIN_RECENT_CODE_SCOPES,
498                    })
499                    .await;
500                Ok(Some(completed))
501            }
502            Err(error) => {
503                let _ = store
504                    .fail_code_index_task(crate::storage::CodeIndexTaskFailure {
505                        task_id: task.task_id,
506                        lease_owner,
507                        attempt_count: task.attempt_count,
508                        error_kind: "code_index".to_owned(),
509                        error_message: error.message.clone(),
510                        retry_backoff_ms: CODE_INDEX_TASK_RETRY_BACKOFF_MS,
511                        max_attempts: CODE_INDEX_TASK_MAX_ATTEMPTS,
512                        now_ms: now_millis(),
513                    })
514                    .await;
515                Err(error)
516            }
517        }
518    }
519
520    /// Recovers code-index worker leases that belonged to exited service processes.
521    pub(crate) async fn recover_orphaned_code_index_tasks_on_startup(
522        &self,
523    ) -> Result<usize, ApiError> {
524        let store = self.store().await.map_err(storage_api_error)?;
525        recover_orphaned_code_index_task_leases(&store, now_millis()).await
526    }
527
528    /// Previews the effective code repository indexing scope without writing rows.
529    pub async fn preview_code_repository_scope(
530        &self,
531        request: CodeIndexRequest,
532        context: RequestContext,
533    ) -> Result<CodeRepositoryScopePreviewResponse, ApiError> {
534        let store = self.store().await.map_err(storage_api_error)?;
535        let status = required_code_repository(&store, &request.repository.repository).await?;
536        let registration = registration_from_status(&status);
537        let selector = request.repository.clone();
538        let preview =
539            run_blocking_code(move || preview_repository_scope(&registration, &selector)).await?;
540        let graph_version = store
541            .current_graph_version()
542            .await
543            .map_err(storage_api_error)?;
544        Ok(CodeRepositoryScopePreviewResponse {
545            metadata: ApiMetadata::graph_only(&context, graph_version),
546            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
547                &status,
548                &request.repository,
549                request.repository.ref_selector.clone(),
550            ),
551            preview,
552        })
553    }
554
555    /// Queries indexed symbols, references, imports, calls, and code chunks.
556    pub async fn query_code_repository(
557        &self,
558        request: CodeRetrievalRequest,
559        context: RequestContext,
560    ) -> Result<CodeRepositoryQueryResponse, ApiError> {
561        let store = self.store().await.map_err(storage_api_error)?;
562        let status = required_code_repository(&store, &request.repository.repository).await?;
563        if request.freshness_policy == FreshnessPolicy::GraphOnly {
564            let graph_version = store
565                .current_graph_version()
566                .await
567                .map_err(storage_api_error)?;
568            let degraded_reason = "graph_only freshness policy selected".to_owned();
569            return Ok(CodeRepositoryQueryResponse {
570                metadata: ApiMetadata::graph_only(&context, graph_version),
571                scope: crate::api::CodeRepositoryScopeMetadata::from_status(
572                    &status,
573                    &request.repository,
574                    request.repository.ref_selector.clone(),
575                ),
576                freshness: crate::api::CodeRepositoryFreshnessDiagnostics::graph_only(
577                    graph_version.get(),
578                    request.freshness_policy,
579                    indexed_source_scope(&status),
580                    request.repository.ref_selector.clone(),
581                    degraded_reason.clone(),
582                ),
583                request,
584                results: Vec::new(),
585                degraded_reason: Some(degraded_reason),
586            });
587        }
588        let requested_ref = request.repository.ref_selector.clone();
589        let mut request = retrieval_request_at_indexed_ref(request, &status).await?;
590        let requested_resolved_ref = request.repository.ref_selector.clone();
591        let freshness_target = request.repository.clone();
592        let mut served_stale_scope = false;
593        let mut stale_reason = None;
594        let scoped_status = match resolved_code_scope_status(&store, &status, &request.repository)
595            .await
596        {
597            Ok(scoped_status) => scoped_status,
598            Err(error) if request.freshness_policy == FreshnessPolicy::AllowStale => {
599                if !active_index_matches_request(&store, &status, &request.repository).await? {
600                    return Err(error);
601                }
602                let Some(stale_status) =
603                    latest_compatible_code_scope_status(&store, &request.repository).await?
604                else {
605                    return Err(error);
606                };
607                let Some(last_indexed_commit) = stale_status.last_indexed_commit.clone() else {
608                    return Err(error);
609                };
610                request.repository.ref_selector = last_indexed_commit;
611                served_stale_scope = true;
612                stale_reason = Some(
613                    "requested ref is not indexed yet; served last completed code index".to_owned(),
614                );
615                stale_status
616            }
617            Err(error) => return Err(error),
618        };
619        if request.freshness_policy == FreshnessPolicy::WaitUntilFresh && scoped_status.stale {
620            return Err(ApiError::invalid_argument(format!(
621                "code repository '{}' scope '{}' is stale; run repo index or repo update before querying with wait_until_fresh",
622                scoped_status.alias,
623                scoped_status
624                    .last_indexed_scope_id
625                    .as_deref()
626                    .unwrap_or("unscoped")
627            )));
628        }
629        let graph_version = store
630            .current_graph_version()
631            .await
632            .map_err(storage_api_error)?;
633        let source_scope = indexed_source_scope(&scoped_status)
634            .ok_or_else(|| missing_indexed_source_scope_error(&scoped_status))?;
635        let mut results = store
636            .search_code_scope(source_scope, request.clone())
637            .await
638            .map_err(storage_api_error)?;
639        let fallback_degraded_reason =
640            apply_code_grep_fallback(&store, &status, &scoped_status, &request, &mut results)
641                .await?;
642        let degraded_reason = results
643            .iter()
644            .find_map(|hit| hit.degraded_reason.clone())
645            .or(fallback_degraded_reason)
646            .or_else(|| scoped_status.degraded_reason.clone())
647            .or_else(|| stale_reason.clone());
648        let mut scope = crate::api::CodeRepositoryScopeMetadata::from_status(
649            &scoped_status,
650            &request.repository,
651            requested_ref.clone(),
652        );
653        if served_stale_scope {
654            scope.stale = true;
655        }
656        let mut metadata = ApiMetadata::graph_only(&context, graph_version);
657        if served_stale_scope {
658            metadata.stale = true;
659        }
660        let freshness = code_query_freshness_diagnostics(
661            &store,
662            CodeQueryFreshnessContext {
663                base_status: &status,
664                scoped_status: &scoped_status,
665                request: &request,
666                requested_ref,
667                requested_resolved_ref,
668                freshness_target,
669                stale_reason,
670                degraded_reason: degraded_reason.clone(),
671                results: &results,
672                graph_version: graph_version.get(),
673            },
674        )
675        .await?;
676        annotate_query_result_staleness(&mut results, &freshness);
677
678        Ok(CodeRepositoryQueryResponse {
679            metadata,
680            scope,
681            freshness,
682            request,
683            results,
684            degraded_reason,
685        })
686    }
687
688    /// Lists configuration-driven feature flags and their code graph relationships.
689    pub async fn query_code_repository_feature_flags(
690        &self,
691        request: CodeFeatureFlagRequest,
692        context: RequestContext,
693    ) -> Result<CodeRepositoryFeatureFlagsResponse, ApiError> {
694        let store = self.store().await.map_err(storage_api_error)?;
695        let status = required_code_repository(&store, &request.repository.repository).await?;
696        if request.freshness_policy == FreshnessPolicy::GraphOnly {
697            let graph_version = store
698                .current_graph_version()
699                .await
700                .map_err(storage_api_error)?;
701            let degraded_reason = "graph_only freshness policy selected".to_owned();
702            return Ok(CodeRepositoryFeatureFlagsResponse {
703                metadata: ApiMetadata::graph_only(&context, graph_version),
704                scope: crate::api::CodeRepositoryScopeMetadata::from_status(
705                    &status,
706                    &request.repository,
707                    request.repository.ref_selector.clone(),
708                ),
709                freshness: crate::api::CodeRepositoryFreshnessDiagnostics::graph_only(
710                    graph_version.get(),
711                    request.freshness_policy,
712                    indexed_source_scope(&status),
713                    request.repository.ref_selector.clone(),
714                    degraded_reason.clone(),
715                ),
716                request,
717                flags: Vec::new(),
718                degraded_reason: Some(degraded_reason),
719            });
720        }
721        let requested_ref = request.repository.ref_selector.clone();
722        let mut request = feature_flag_request_at_indexed_ref(request, &status).await?;
723        let requested_resolved_ref = request.repository.ref_selector.clone();
724        let freshness_target = request.repository.clone();
725        let mut served_stale_scope = false;
726        let mut stale_reason = None;
727        let scoped_status = match resolved_code_scope_status(&store, &status, &request.repository)
728            .await
729        {
730            Ok(scoped_status) => scoped_status,
731            Err(error) if request.freshness_policy == FreshnessPolicy::AllowStale => {
732                if !active_index_matches_request(&store, &status, &request.repository).await? {
733                    return Err(error);
734                }
735                let Some(stale_status) =
736                    latest_compatible_code_scope_status(&store, &request.repository).await?
737                else {
738                    return Err(error);
739                };
740                let Some(last_indexed_commit) = stale_status.last_indexed_commit.clone() else {
741                    return Err(error);
742                };
743                request.repository.ref_selector = last_indexed_commit;
744                served_stale_scope = true;
745                stale_reason = Some(
746                    "requested ref is not indexed yet; served last completed code index".to_owned(),
747                );
748                stale_status
749            }
750            Err(error) => return Err(error),
751        };
752        if request.freshness_policy == FreshnessPolicy::WaitUntilFresh && scoped_status.stale {
753            return Err(ApiError::invalid_argument(format!(
754                "code repository '{}' scope '{}' is stale; run repo index or repo update before querying feature flags with wait_until_fresh",
755                scoped_status.alias,
756                scoped_status
757                    .last_indexed_scope_id
758                    .as_deref()
759                    .unwrap_or("unscoped")
760            )));
761        }
762        let graph_version = store
763            .current_graph_version()
764            .await
765            .map_err(storage_api_error)?;
766        let source_scope = indexed_source_scope(&scoped_status)
767            .ok_or_else(|| missing_indexed_source_scope_error(&scoped_status))?;
768        let flags = store
769            .search_code_feature_flags_scope(source_scope, request.clone())
770            .await
771            .map_err(storage_api_error)?;
772        let mut scope = crate::api::CodeRepositoryScopeMetadata::from_status(
773            &scoped_status,
774            &request.repository,
775            requested_ref.clone(),
776        );
777        if served_stale_scope {
778            scope.stale = true;
779        }
780        let degraded_reason = scoped_status
781            .degraded_reason
782            .clone()
783            .or_else(|| stale_reason.clone());
784        let mut metadata = ApiMetadata::graph_only(&context, graph_version);
785        if served_stale_scope {
786            metadata.stale = true;
787        }
788        let freshness = code_feature_flag_freshness_diagnostics(
789            &store,
790            CodeFeatureFlagFreshnessContext {
791                base_status: &status,
792                scoped_status: &scoped_status,
793                request: &request,
794                requested_ref,
795                requested_resolved_ref,
796                freshness_target,
797                stale_reason,
798                degraded_reason: degraded_reason.clone(),
799                flags: &flags,
800                graph_version: graph_version.get(),
801            },
802        )
803        .await?;
804
805        Ok(CodeRepositoryFeatureFlagsResponse {
806            metadata,
807            scope,
808            freshness,
809            request,
810            flags,
811            degraded_reason,
812        })
813    }
814
815    /// Returns impact radius for a Git diff using the indexed code graph.
816    pub async fn impact_code_repository(
817        &self,
818        mut request: CodeImpactRequest,
819        context: RequestContext,
820    ) -> Result<CodeRepositoryImpactResponse, ApiError> {
821        let store = self.store().await.map_err(storage_api_error)?;
822        let status = required_code_repository(&store, &request.repository.repository).await?;
823        let head_commit =
824            resolve_code_ref_for_selector(&status, &request.repository, request.head_ref.clone())
825                .await?;
826        request.repository.ref_selector = head_commit.clone();
827        let scoped_status =
828            resolved_code_scope_status(&store, &status, &request.repository).await?;
829        let root = PathBuf::from(status.root_path.clone());
830        let base_ref = request.base_ref.clone();
831        let head_ref = head_commit.clone();
832        let path_filters = scoped_status.path_filters.clone();
833        let language_filters = scoped_status.language_filters.clone();
834        let base_fingerprints = if base_ref.starts_with("filesystem:") {
835            let mut base_selector = request.repository.clone();
836            base_selector.ref_selector = base_ref.clone();
837            match resolved_code_scope_status(&store, &status, &base_selector).await {
838                Ok(base_status) => match base_status.last_indexed_scope_id {
839                    Some(source_scope) => Some(
840                        store
841                            .code_file_fingerprints_for_scope(source_scope)
842                            .await
843                            .map_err(storage_api_error)?,
844                    ),
845                    None => None,
846                },
847                Err(_) => None,
848            }
849        } else {
850            None
851        };
852        let changed_paths = if let Some(base_fingerprints) = base_fingerprints {
853            run_blocking_code(move || {
854                let previous_hashes = base_fingerprints
855                    .into_iter()
856                    .map(|file| (file.path, file.blob_hash))
857                    .collect::<BTreeMap<_, _>>();
858                changed_paths_for_filesystem_diff(
859                    &root,
860                    &head_ref,
861                    &path_filters,
862                    &language_filters,
863                    &previous_hashes,
864                )
865            })
866            .await?
867        } else {
868            run_blocking_code(move || {
869                changed_paths_for_diff_with_filters(
870                    root,
871                    &base_ref,
872                    &head_ref,
873                    &path_filters,
874                    &language_filters,
875                )
876            })
877            .await?
878        };
879        let registration = registration_from_status(&status);
880        let path_groups = {
881            let registration = registration.clone();
882            let selector = request.repository.clone();
883            let changed_paths = changed_paths.clone();
884            run_blocking_code(move || {
885                partition_changed_paths_for_selector(&registration, &selector, changed_paths)
886            })
887            .await?
888        };
889        let selector = request.repository.clone();
890        let base_ref = request.base_ref.clone();
891        let head_ref = head_commit;
892        let deleted_symbol_names = run_blocking_code(move || {
893            deleted_symbol_names_for_diff(&registration, &selector, &base_ref, &head_ref)
894        })
895        .await?;
896        let source_scope = indexed_source_scope(&scoped_status)
897            .ok_or_else(|| missing_indexed_source_scope_error(&scoped_status))?;
898        let results = store
899            .analyze_code_impact_scope(
900                source_scope,
901                request.clone(),
902                CodeImpactChanges {
903                    paths: changed_paths.clone(),
904                    deleted_symbol_names,
905                },
906            )
907            .await
908            .map_err(storage_api_error)?;
909        let graph_version = store
910            .current_graph_version()
911            .await
912            .map_err(storage_api_error)?;
913        let scope = crate::api::CodeRepositoryScopeMetadata::from_status(
914            &scoped_status,
915            &request.repository,
916            request.head_ref.clone(),
917        );
918
919        Ok(CodeRepositoryImpactResponse {
920            metadata: ApiMetadata::graph_only(&context, graph_version),
921            scope,
922            request,
923            path_groups,
924            results,
925        })
926    }
927
928    pub async fn code_repository_status(
929        &self,
930        selector: CodeRepositorySelector,
931        context: RequestContext,
932    ) -> Result<CodeRepositoryStatusResponse, ApiError> {
933        let store = self.store().await.map_err(storage_api_error)?;
934        let status = required_code_repository(&store, &selector.repository).await?;
935        recover_code_index_task_leases(&store, now_millis()).await?;
936        let active_task = store
937            .active_code_index_task(status.repository_id.clone())
938            .await
939            .map_err(storage_api_error)?;
940        let checkpoint = code_status_checkpoint(&store, &status, active_task.as_ref()).await?;
941        let retention = store
942            .code_scope_retention(status.repository_id.clone())
943            .await
944            .map_err(storage_api_error)?;
945        let graph_version = store
946            .current_graph_version()
947            .await
948            .map_err(storage_api_error)?;
949
950        Ok(CodeRepositoryStatusResponse {
951            metadata: ApiMetadata::graph_only(&context, graph_version),
952            status,
953            active_task,
954            checkpoint,
955            retention,
956        })
957    }
958
959    pub(crate) async fn code_repository_is_registered(
960        &self,
961        repository: String,
962    ) -> Result<bool, ApiError> {
963        let selector = CodeRepositorySelector::new(repository, "HEAD", Vec::new(), Vec::new())
964            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
965        let store = self.store().await.map_err(storage_api_error)?;
966        store
967            .code_repository_status(selector.repository)
968            .await
969            .map(|status| status.is_some())
970            .map_err(storage_api_error)
971    }
972
973    /// Builds a reusable operations report for a registered code repository.
974    pub async fn code_repository_report(
975        &self,
976        selector: CodeRepositorySelector,
977        context: RequestContext,
978    ) -> Result<CodeRepositoryReportResponse, ApiError> {
979        let store = self.store().await.map_err(storage_api_error)?;
980        let status = required_code_repository(&store, &selector.repository).await?;
981        let report = store
982            .code_repository_report(status.repository_id.clone())
983            .await
984            .map_err(storage_api_error)?;
985        let graph_version = store
986            .current_graph_version()
987            .await
988            .map_err(storage_api_error)?;
989
990        Ok(CodeRepositoryReportResponse {
991            metadata: ApiMetadata::graph_only(&context, graph_version),
992            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
993                &status,
994                &selector,
995                selector.ref_selector.clone(),
996            ),
997            report,
998        })
999    }
1000}