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