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