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