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