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