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