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