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