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