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        CodeRepositoryQueryResponse, CodeRepositoryRegisterRequest, CodeRepositoryRegisterResponse,
7        CodeRepositoryReportResponse, CodeRepositoryScopePreviewResponse,
8        CodeRepositoryStatusResponse, RequestContext,
9    },
10    code::{
11        CodeIndexError, build_index_snapshot, changed_paths_for_diff,
12        deleted_symbol_names_for_diff, partition_changed_paths_for_selector,
13        prepare_full_index_plan, preview_repository_scope, register_repository,
14        resolve_repository_ref, resolve_repository_snapshot,
15    },
16    domain::{
17        CodeImpactRequest, CodeIndexMode, CodeIndexRequest, CodeIndexResourceBudget,
18        CodeRepositoryRegistration, CodeRepositorySelector, CodeRepositoryStatus,
19        CodeRetrievalRequest, FreshnessPolicy,
20    },
21    storage::{CodeImpactChanges, StorageError},
22};
23
24use super::RelayKnowledgeService;
25
26impl RelayKnowledgeService {
27    /// Registers a Git repository as a code source.
28    pub async fn register_code_repository(
29        &self,
30        request: CodeRepositoryRegisterRequest,
31        context: RequestContext,
32    ) -> Result<CodeRepositoryRegisterResponse, ApiError> {
33        let registration = run_blocking_code(move || {
34            register_repository(
35                request.root_path,
36                request.alias,
37                request.path_filters,
38                request.language_filters,
39            )
40        })
41        .await?;
42        let store = self.store().await.map_err(storage_api_error)?;
43        let status = store
44            .upsert_code_repository(registration.clone())
45            .await
46            .map_err(storage_api_error)?;
47        let graph_version = store
48            .current_graph_version()
49            .await
50            .map_err(storage_api_error)?;
51
52        Ok(CodeRepositoryRegisterResponse {
53            metadata: ApiMetadata::graph_only(&context, graph_version),
54            registration,
55            status,
56        })
57    }
58
59    /// Builds or updates the tree-sitter code index for a registered repository.
60    pub async fn index_code_repository(
61        &self,
62        request: CodeIndexRequest,
63        context: RequestContext,
64    ) -> Result<CodeRepositoryIndexResponse, ApiError> {
65        let store = self.store().await.map_err(storage_api_error)?;
66        let status = required_code_repository(&store, &request.repository.repository).await?;
67        if let Some(response) = self
68            .fresh_full_index_response(&store, &status, &request, &context)
69            .await?
70        {
71            return Ok(response);
72        }
73        let registration = registration_from_status(&status);
74        let selector = request.repository.clone();
75        let summary = if request.mode == CodeIndexMode::Full {
76            let resource_budget = CodeIndexResourceBudget::default();
77            let mut plan = run_blocking_code(move || {
78                prepare_full_index_plan(registration, selector, resource_budget)
79            })
80            .await?;
81            let session = plan.session();
82            store
83                .begin_code_index_session(session.clone())
84                .await
85                .map_err(storage_api_error)?;
86            loop {
87                let (next_plan, batch) = run_blocking_code(move || plan.parse_next_batch()).await?;
88                plan = next_plan;
89                let Some(batch) = batch else {
90                    break;
91                };
92                store
93                    .apply_code_index_batch(batch)
94                    .await
95                    .map_err(storage_api_error)?;
96            }
97            store
98                .finalize_code_index_session(session)
99                .await
100                .map_err(storage_api_error)?
101        } else {
102            let previous = previous_fingerprints_for_index(&store, &status, &request).await?;
103            let mode = request.mode;
104            let snapshot = run_blocking_code(move || {
105                build_index_snapshot(&registration, &selector, mode, previous)
106            })
107            .await?;
108            store
109                .apply_code_index_snapshot(snapshot)
110                .await
111                .map_err(storage_api_error)?
112        };
113        let status = store
114            .code_repository_status(summary.repository_id.clone())
115            .await
116            .map_err(storage_api_error)?
117            .ok_or_else(|| ApiError::storage_unavailable("code repository status is missing"))?;
118        let graph_version = store
119            .current_graph_version()
120            .await
121            .map_err(storage_api_error)?;
122
123        Ok(CodeRepositoryIndexResponse {
124            metadata: ApiMetadata::graph_only(&context, graph_version),
125            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
126                &status,
127                &request.repository,
128                request.repository.ref_selector.clone(),
129            ),
130            summary,
131            status,
132        })
133    }
134
135    async fn fresh_full_index_response(
136        &self,
137        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
138        status: &CodeRepositoryStatus,
139        request: &CodeIndexRequest,
140        context: &RequestContext,
141    ) -> Result<Option<CodeRepositoryIndexResponse>, ApiError> {
142        if request.mode != CodeIndexMode::Full {
143            return Ok(None);
144        }
145        let registration = registration_from_status(status);
146        let selector = request.repository.clone();
147        let (resolved_commit_sha, tree_hash) = run_blocking_code(move || {
148            resolve_repository_snapshot(&registration.root_path, &selector.ref_selector)
149        })
150        .await?;
151        let path_filters = merged_filters(&status.path_filters, &request.repository.path_filters);
152        let language_filters = merged_filters(
153            &status.language_filters,
154            &request.repository.language_filters,
155        );
156        let scoped_status = store
157            .code_repository_scope_status(
158                request.repository.repository.clone(),
159                resolved_commit_sha.clone(),
160                path_filters,
161                language_filters,
162            )
163            .await
164            .map_err(storage_api_error)?;
165        let Some(scoped_status) = scoped_status else {
166            return Ok(None);
167        };
168        if scoped_status.stale || scoped_status.tree_hash.as_deref() != Some(tree_hash.as_str()) {
169            return Ok(None);
170        }
171        let graph_version = store
172            .current_graph_version()
173            .await
174            .map_err(storage_api_error)?;
175        let report = store
176            .code_repository_report(scoped_status.repository_id.clone())
177            .await
178            .map_err(storage_api_error)?;
179        let summary = crate::domain::CodeIndexSummary {
180            repository_id: scoped_status.repository_id.clone(),
181            source_scope: scoped_status
182                .last_indexed_scope_id
183                .clone()
184                .unwrap_or_default(),
185            resolved_commit_sha,
186            tree_hash,
187            indexed_file_count: scoped_status.indexed_file_count,
188            changed_path_count: 0,
189            skipped_unchanged_count: scoped_status.indexed_file_count,
190            deleted_path_count: 0,
191            symbol_count: scoped_status.symbol_count,
192            reference_count: scoped_status.reference_count,
193            chunk_count: scoped_status.chunk_count,
194            degraded_file_count: report.degraded_file_count,
195            progress: crate::domain::CodeIndexProgressSummary {
196                git_file_count: scoped_status.indexed_file_count,
197                blob_read_count: 0,
198                parsed_file_count: 0,
199                sqlite_write_count: 0,
200                skipped_file_count: scoped_status.indexed_file_count,
201                degraded_file_count: report.degraded_file_count,
202                batch_count: 0,
203                checkpoint_file_count: scoped_status.indexed_file_count,
204                resource_budget: crate::domain::CodeIndexResourceBudget::default(),
205            },
206        };
207
208        Ok(Some(CodeRepositoryIndexResponse {
209            metadata: ApiMetadata::graph_only(context, graph_version),
210            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
211                &scoped_status,
212                &request.repository,
213                request.repository.ref_selector.clone(),
214            ),
215            summary,
216            status: scoped_status,
217        }))
218    }
219
220    /// Previews the effective code repository indexing scope without writing rows.
221    pub async fn preview_code_repository_scope(
222        &self,
223        request: CodeIndexRequest,
224        context: RequestContext,
225    ) -> Result<CodeRepositoryScopePreviewResponse, ApiError> {
226        let store = self.store().await.map_err(storage_api_error)?;
227        let status = required_code_repository(&store, &request.repository.repository).await?;
228        let registration = registration_from_status(&status);
229        let selector = request.repository.clone();
230        let preview =
231            run_blocking_code(move || preview_repository_scope(&registration, &selector)).await?;
232        let graph_version = store
233            .current_graph_version()
234            .await
235            .map_err(storage_api_error)?;
236        Ok(CodeRepositoryScopePreviewResponse {
237            metadata: ApiMetadata::graph_only(&context, graph_version),
238            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
239                &status,
240                &request.repository,
241                request.repository.ref_selector.clone(),
242            ),
243            preview,
244        })
245    }
246
247    /// Queries indexed symbols, references, imports, calls, and code chunks.
248    pub async fn query_code_repository(
249        &self,
250        request: CodeRetrievalRequest,
251        context: RequestContext,
252    ) -> Result<CodeRepositoryQueryResponse, ApiError> {
253        let store = self.store().await.map_err(storage_api_error)?;
254        let status = required_code_repository(&store, &request.repository.repository).await?;
255        if request.freshness_policy == FreshnessPolicy::GraphOnly {
256            let graph_version = store
257                .current_graph_version()
258                .await
259                .map_err(storage_api_error)?;
260            return Ok(CodeRepositoryQueryResponse {
261                metadata: ApiMetadata::graph_only(&context, graph_version),
262                scope: crate::api::CodeRepositoryScopeMetadata::from_status(
263                    &status,
264                    &request.repository,
265                    request.repository.ref_selector.clone(),
266                ),
267                request,
268                results: Vec::new(),
269                degraded_reason: Some("graph_only freshness policy selected".to_owned()),
270            });
271        }
272        let requested_ref = request.repository.ref_selector.clone();
273        let request = retrieval_request_at_indexed_ref(request, &status).await?;
274        let scoped_status =
275            resolved_code_scope_status(&store, &status, &request.repository).await?;
276        if request.freshness_policy == FreshnessPolicy::WaitUntilFresh && scoped_status.stale {
277            return Err(ApiError::invalid_argument(format!(
278                "code repository '{}' scope '{}' is stale; run repo index or repo update before querying with wait_until_fresh",
279                scoped_status.alias,
280                scoped_status
281                    .last_indexed_scope_id
282                    .as_deref()
283                    .unwrap_or("unscoped")
284            )));
285        }
286        let graph_version = store
287            .current_graph_version()
288            .await
289            .map_err(storage_api_error)?;
290        let results = store
291            .search_code(request.clone())
292            .await
293            .map_err(storage_api_error)?;
294        let degraded_reason = results
295            .iter()
296            .find_map(|hit| hit.degraded_reason.clone())
297            .or_else(|| scoped_status.degraded_reason.clone());
298        let scope = crate::api::CodeRepositoryScopeMetadata::from_status(
299            &scoped_status,
300            &request.repository,
301            requested_ref,
302        );
303
304        Ok(CodeRepositoryQueryResponse {
305            metadata: ApiMetadata::graph_only(&context, graph_version),
306            scope,
307            request,
308            results,
309            degraded_reason,
310        })
311    }
312
313    /// Returns impact radius for a Git diff using the indexed code graph.
314    pub async fn impact_code_repository(
315        &self,
316        mut request: CodeImpactRequest,
317        context: RequestContext,
318    ) -> Result<CodeRepositoryImpactResponse, ApiError> {
319        let store = self.store().await.map_err(storage_api_error)?;
320        let status = required_code_repository(&store, &request.repository.repository).await?;
321        let head_commit = resolve_code_ref(&status, request.head_ref.clone()).await?;
322        request.repository.ref_selector = head_commit.clone();
323        let scoped_status =
324            resolved_code_scope_status(&store, &status, &request.repository).await?;
325        let root = PathBuf::from(status.root_path.clone());
326        let base_ref = request.base_ref.clone();
327        let head_ref = head_commit.clone();
328        let changed_paths =
329            run_blocking_code(move || changed_paths_for_diff(root, &base_ref, &head_ref)).await?;
330        let registration = registration_from_status(&status);
331        let path_groups = {
332            let registration = registration.clone();
333            let selector = request.repository.clone();
334            let changed_paths = changed_paths.clone();
335            run_blocking_code(move || {
336                partition_changed_paths_for_selector(&registration, &selector, changed_paths)
337            })
338            .await?
339        };
340        let selector = request.repository.clone();
341        let base_ref = request.base_ref.clone();
342        let head_ref = head_commit;
343        let deleted_symbol_names = run_blocking_code(move || {
344            deleted_symbol_names_for_diff(&registration, &selector, &base_ref, &head_ref)
345        })
346        .await?;
347        let results = store
348            .analyze_code_impact(
349                request.clone(),
350                CodeImpactChanges {
351                    paths: changed_paths.clone(),
352                    deleted_symbol_names,
353                },
354            )
355            .await
356            .map_err(storage_api_error)?;
357        let graph_version = store
358            .current_graph_version()
359            .await
360            .map_err(storage_api_error)?;
361        let scope = crate::api::CodeRepositoryScopeMetadata::from_status(
362            &scoped_status,
363            &request.repository,
364            request.head_ref.clone(),
365        );
366
367        Ok(CodeRepositoryImpactResponse {
368            metadata: ApiMetadata::graph_only(&context, graph_version),
369            scope,
370            request,
371            path_groups,
372            results,
373        })
374    }
375
376    /// Returns the current code repository index status.
377    pub async fn code_repository_status(
378        &self,
379        selector: CodeRepositorySelector,
380        context: RequestContext,
381    ) -> Result<CodeRepositoryStatusResponse, ApiError> {
382        let store = self.store().await.map_err(storage_api_error)?;
383        let status = required_code_repository(&store, &selector.repository).await?;
384        let graph_version = store
385            .current_graph_version()
386            .await
387            .map_err(storage_api_error)?;
388
389        Ok(CodeRepositoryStatusResponse {
390            metadata: ApiMetadata::graph_only(&context, graph_version),
391            status,
392        })
393    }
394
395    /// Checks whether a repository selector resolves to a registered code source.
396    pub(crate) async fn code_repository_is_registered(
397        &self,
398        repository: String,
399    ) -> Result<bool, ApiError> {
400        let selector = CodeRepositorySelector::new(repository, "HEAD", Vec::new(), Vec::new())
401            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
402        let store = self.store().await.map_err(storage_api_error)?;
403        store
404            .code_repository_status(selector.repository)
405            .await
406            .map(|status| status.is_some())
407            .map_err(storage_api_error)
408    }
409
410    /// Builds a reusable operations report for a registered code repository.
411    pub async fn code_repository_report(
412        &self,
413        selector: CodeRepositorySelector,
414        context: RequestContext,
415    ) -> Result<CodeRepositoryReportResponse, ApiError> {
416        let store = self.store().await.map_err(storage_api_error)?;
417        let status = required_code_repository(&store, &selector.repository).await?;
418        let report = store
419            .code_repository_report(status.repository_id.clone())
420            .await
421            .map_err(storage_api_error)?;
422        let graph_version = store
423            .current_graph_version()
424            .await
425            .map_err(storage_api_error)?;
426
427        Ok(CodeRepositoryReportResponse {
428            metadata: ApiMetadata::graph_only(&context, graph_version),
429            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
430                &status,
431                &selector,
432                selector.ref_selector.clone(),
433            ),
434            report,
435        })
436    }
437}
438
439async fn required_code_repository(
440    store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
441    repository: &str,
442) -> Result<crate::domain::CodeRepositoryStatus, ApiError> {
443    store
444        .code_repository_status(repository.to_owned())
445        .await
446        .map_err(storage_api_error)?
447        .ok_or_else(|| {
448            ApiError::invalid_argument(format!("code repository '{repository}' is not registered"))
449        })
450}
451
452fn registration_from_status(
453    status: &crate::domain::CodeRepositoryStatus,
454) -> CodeRepositoryRegistration {
455    CodeRepositoryRegistration {
456        repository_id: status.repository_id.clone(),
457        alias: status.alias.clone(),
458        root_path: status.root_path.clone(),
459        path_filters: status.path_filters.clone(),
460        language_filters: status.language_filters.clone(),
461    }
462}
463
464async fn previous_fingerprints_for_index(
465    store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
466    status: &CodeRepositoryStatus,
467    request: &CodeIndexRequest,
468) -> Result<Vec<crate::domain::CodeFileFingerprint>, ApiError> {
469    let CodeIndexMode::Incremental { base_ref, .. } = &request.mode else {
470        return store
471            .code_file_fingerprints(status.repository_id.clone())
472            .await
473            .map_err(storage_api_error);
474    };
475    let base_commit = resolve_code_ref(status, base_ref.clone()).await?;
476    let path_filters = merged_filters(&status.path_filters, &request.repository.path_filters);
477    let language_filters = merged_filters(
478        &status.language_filters,
479        &request.repository.language_filters,
480    );
481    let base_scope = store
482        .code_repository_scope_status(
483            request.repository.repository.clone(),
484            base_commit.clone(),
485            path_filters,
486            language_filters,
487        )
488        .await
489        .map_err(storage_api_error)?
490        .ok_or_else(|| {
491            ApiError::invalid_argument(format!(
492                "incremental base ref '{}' resolves to {}, but code repository '{}' has no matching indexed base scope; run repo index --ref {} before repo update",
493                base_ref, base_commit, status.alias, base_ref
494            ))
495        })?;
496    if base_scope.stale {
497        return Err(ApiError::invalid_argument(format!(
498            "incremental base ref '{}' resolves to a stale indexed scope {}; refresh or reindex the base before repo update",
499            base_ref,
500            base_scope
501                .last_indexed_scope_id
502                .as_deref()
503                .unwrap_or("unscoped")
504        )));
505    }
506    let source_scope = base_scope.last_indexed_scope_id.ok_or_else(|| {
507        ApiError::invalid_argument(format!(
508            "incremental base ref '{}' has no persisted source scope",
509            base_ref
510        ))
511    })?;
512
513    store
514        .code_file_fingerprints_for_scope(source_scope)
515        .await
516        .map_err(storage_api_error)
517}
518
519async fn retrieval_request_at_indexed_ref(
520    mut request: CodeRetrievalRequest,
521    status: &CodeRepositoryStatus,
522) -> Result<CodeRetrievalRequest, ApiError> {
523    request.repository.ref_selector =
524        indexed_commit_for_ref(status, request.repository.ref_selector.clone()).await?;
525
526    Ok(request)
527}
528
529async fn resolved_code_scope_status(
530    store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
531    status: &CodeRepositoryStatus,
532    selector: &CodeRepositorySelector,
533) -> Result<CodeRepositoryStatus, ApiError> {
534    let path_filters = merged_filters(&status.path_filters, &selector.path_filters);
535    let language_filters = merged_filters(&status.language_filters, &selector.language_filters);
536    let exact_scope = store
537        .code_repository_scope_status(
538            selector.repository.clone(),
539            selector.ref_selector.clone(),
540            path_filters,
541            language_filters,
542        )
543        .await
544        .map_err(storage_api_error)?;
545    let scoped_status = match exact_scope {
546        Some(status) => Some(status),
547        None if (!selector.path_filters.is_empty() || !selector.language_filters.is_empty())
548            && selector_filters_fit_indexed_scope(status, selector) =>
549        {
550            store
551                .code_repository_scope_status(
552                    selector.repository.clone(),
553                    selector.ref_selector.clone(),
554                    status.path_filters.clone(),
555                    status.language_filters.clone(),
556                )
557                .await
558                .map_err(storage_api_error)?
559        }
560        None => None,
561    };
562    scoped_status.ok_or_else(|| {
563        ApiError::invalid_argument(format!(
564            "code repository '{}' has no index for ref {} and requested filters",
565            selector.repository, selector.ref_selector
566        ))
567    })
568}
569
570fn merged_filters(left: &[String], right: &[String]) -> Vec<String> {
571    let mut merged = Vec::new();
572    for value in left.iter().chain(right.iter()) {
573        if !merged.contains(value) {
574            merged.push(value.clone());
575        }
576    }
577
578    merged
579}
580
581fn selector_filters_fit_indexed_scope(
582    status: &CodeRepositoryStatus,
583    selector: &CodeRepositorySelector,
584) -> bool {
585    requested_paths_fit_indexed_scope(&status.path_filters, &selector.path_filters)
586        && requested_languages_fit_indexed_scope(
587            &status.language_filters,
588            &selector.language_filters,
589        )
590}
591
592fn requested_paths_fit_indexed_scope(
593    indexed_filters: &[String],
594    selector_filters: &[String],
595) -> bool {
596    selector_filters.is_empty()
597        || indexed_filters.is_empty()
598        || selector_filters.iter().all(|selector_filter| {
599            indexed_filters
600                .iter()
601                .any(|indexed_filter| path_filter_covers(indexed_filter, selector_filter))
602        })
603}
604
605fn requested_languages_fit_indexed_scope(
606    indexed_filters: &[String],
607    selector_filters: &[String],
608) -> bool {
609    selector_filters.is_empty()
610        || indexed_filters.is_empty()
611        || selector_filters
612            .iter()
613            .all(|selector_filter| indexed_filters.contains(selector_filter))
614}
615
616fn path_filter_covers(indexed_filter: &str, selector_filter: &str) -> bool {
617    let indexed_filter = normalize_path_filter(indexed_filter);
618    let selector_filter = normalize_path_filter(selector_filter);
619    indexed_filter == "."
620        || (!indexed_filter.is_empty()
621            && !selector_filter.is_empty()
622            && (selector_filter == indexed_filter
623                || selector_filter.starts_with(&format!("{indexed_filter}/"))))
624}
625
626fn normalize_path_filter(filter: &str) -> &str {
627    let mut filter = filter.trim_end_matches(['/', '\\']);
628    while let Some(stripped) = filter.strip_prefix("./") {
629        filter = stripped;
630    }
631
632    filter
633}
634
635async fn indexed_commit_for_ref(
636    status: &CodeRepositoryStatus,
637    ref_selector: String,
638) -> Result<String, ApiError> {
639    if ref_selector == "worktree" {
640        if is_worktree_overlay(status) {
641            return status.last_indexed_commit.clone().ok_or_else(|| {
642                ApiError::invalid_argument(format!(
643                    "code repository '{}' has no active worktree overlay",
644                    status.alias
645                ))
646            });
647        }
648        return Err(ApiError::invalid_argument(format!(
649            "code repository '{}' has no active worktree overlay",
650            status.alias
651        )));
652    }
653
654    resolve_code_ref(status, ref_selector).await
655}
656
657fn is_worktree_overlay(status: &CodeRepositoryStatus) -> bool {
658    status
659        .last_indexed_commit
660        .as_deref()
661        .is_some_and(|value| value.starts_with("worktree:"))
662        || status
663            .tree_hash
664            .as_deref()
665            .is_some_and(|value| value.starts_with("worktree:"))
666}
667
668async fn resolve_code_ref(
669    status: &CodeRepositoryStatus,
670    ref_selector: String,
671) -> Result<String, ApiError> {
672    let root = PathBuf::from(status.root_path.clone());
673
674    run_blocking_code(move || resolve_repository_ref(root, &ref_selector)).await
675}
676
677async fn run_blocking_code<T, F>(operation: F) -> Result<T, ApiError>
678where
679    T: Send + 'static,
680    F: FnOnce() -> Result<T, CodeIndexError> + Send + 'static,
681{
682    tokio::task::spawn_blocking(operation)
683        .await
684        .map_err(|error| ApiError::storage_unavailable(error.to_string()))?
685        .map_err(code_api_error)
686}
687
688fn code_api_error(error: CodeIndexError) -> ApiError {
689    match error {
690        CodeIndexError::InvalidInput(message) => ApiError::invalid_argument(message),
691        CodeIndexError::Git { .. } | CodeIndexError::Io(_) | CodeIndexError::TreeSitter(_) => {
692            ApiError::storage_unavailable(error.to_string())
693        }
694    }
695}
696
697fn storage_api_error(error: StorageError) -> ApiError {
698    ApiError::storage_unavailable(error.to_string())
699}