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