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