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 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 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(®istration, &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(®istration.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 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(®istration, &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 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 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(®istration, &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(®istration, &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 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 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 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}