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