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