Skip to main content

relay_knowledge/application/
code_service.rs

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