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