diff --git a/java-ecosystem/libs/analysis-engine/src/main/java/org/rostilos/codecrow/analysisengine/processor/analysis/PullRequestAnalysisProcessor.java b/java-ecosystem/libs/analysis-engine/src/main/java/org/rostilos/codecrow/analysisengine/processor/analysis/PullRequestAnalysisProcessor.java index a8a83947..4b458ecd 100644 --- a/java-ecosystem/libs/analysis-engine/src/main/java/org/rostilos/codecrow/analysisengine/processor/analysis/PullRequestAnalysisProcessor.java +++ b/java-ecosystem/libs/analysis-engine/src/main/java/org/rostilos/codecrow/analysisengine/processor/analysis/PullRequestAnalysisProcessor.java @@ -22,7 +22,6 @@ import org.rostilos.codecrow.analysisengine.exception.AnalysisLockedException; import org.rostilos.codecrow.analysisengine.service.AnalysisLockService; import org.rostilos.codecrow.analysisengine.service.PullRequestService; -import org.rostilos.codecrow.analysisapi.rag.RagOperationsService; import org.rostilos.codecrow.commitgraph.service.AnalyzedCommitService; import org.rostilos.codecrow.analysisengine.service.vcs.VcsAiClientService; import org.rostilos.codecrow.analysisengine.service.vcs.VcsReportingService; @@ -72,7 +71,6 @@ public class PullRequestAnalysisProcessor { private final AiAnalysisClient aiAnalysisClient; private final VcsServiceFactory vcsServiceFactory; private final AnalysisLockService analysisLockService; - private final RagOperationsService ragOperationsService; private final ApplicationEventPublisher eventPublisher; private final AnalyzedCommitService analyzedCommitService; private final VcsClientProvider vcsClientProvider; @@ -98,7 +96,6 @@ public PullRequestAnalysisProcessor( FileSnapshotService fileSnapshotService, PrIssueTrackingService prIssueTrackingService, AstScopeEnricher astScopeEnricher, - @Autowired(required = false) RagOperationsService ragOperationsService, @Autowired(required = false) ApplicationEventPublisher eventPublisher ) { this.codeAnalysisService = codeAnalysisService; @@ -107,7 +104,6 @@ public PullRequestAnalysisProcessor( this.aiAnalysisClient = aiAnalysisClient; this.vcsServiceFactory = vcsServiceFactory; this.analysisLockService = analysisLockService; - this.ragOperationsService = ragOperationsService; this.eventPublisher = eventPublisher; this.analyzedCommitService = analyzedCommitService; this.vcsClientProvider = vcsClientProvider; @@ -246,9 +242,6 @@ public Map process( project.getId(), request.getPullRequestId()); } - // Only prepare RAG after the exact snapshot/configuration cache has missed. - ensureRagIndexForTargetBranch(project, request.getTargetBranchName(), consumer); - Map aiResponse = aiAnalysisClient.performAnalysis(aiRequest, event -> { log.debug("Received event from AI client: type={}", event.get("type")); emitEvent(consumer, event); @@ -944,41 +937,6 @@ static List selectPrEvidenceCommits( return selected; } - /** - * Ensures RAG index is up-to-date for the PR target branch. - *

- * For PRs targeting the main branch: - * - Checks if the main RAG index commit matches the current target branch HEAD - * - If outdated, performs incremental update before analysis - *

- * For PRs targeting non-main branches with multi-branch enabled: - * - First ensures the main index is up to date - * - Then ensures branch index exists and is up to date for the target branch - *

- * This ensures analysis always uses the most current codebase context. - */ - private void ensureRagIndexForTargetBranch(Project project, String targetBranch, EventConsumer consumer) { - if (ragOperationsService == null) { - log.debug("RagOperationsService not available - skipping RAG index check for target branch"); - return; - } - - try { - boolean ready = ragOperationsService.ensureRagIndexUpToDate( - project, - targetBranch, - event -> emitEvent(consumer, event)); - if (ready) { - log.info("RAG index ensured up-to-date for PR target branch: project={}, branch={}", - project.getId(), targetBranch); - } - } catch (Exception e) { - log.warn( - "Failed to ensure RAG index up-to-date for target branch (non-critical): project={}, branch={}, error={}", - project.getId(), targetBranch, e.getMessage()); - } - } - /** * Publishes an AnalysisStartedEvent for PR analysis. */ diff --git a/java-ecosystem/libs/analysis-engine/src/main/java/org/rostilos/codecrow/analysisengine/service/branch/BranchAnalysisGateService.java b/java-ecosystem/libs/analysis-engine/src/main/java/org/rostilos/codecrow/analysisengine/service/branch/BranchAnalysisGateService.java index 6718f620..08d5069e 100644 --- a/java-ecosystem/libs/analysis-engine/src/main/java/org/rostilos/codecrow/analysisengine/service/branch/BranchAnalysisGateService.java +++ b/java-ecosystem/libs/analysis-engine/src/main/java/org/rostilos/codecrow/analysisengine/service/branch/BranchAnalysisGateService.java @@ -29,6 +29,7 @@ public class BranchAnalysisGateService { private static final Logger log = LoggerFactory.getLogger(BranchAnalysisGateService.class); + private static final long WAIT_STATUS_INTERVAL_NANOS = TimeUnit.SECONDS.toNanos(60); private final JobRepository jobRepository; @@ -90,6 +91,7 @@ public GateResult awaitTurn( Consumer> consumer) { long timeoutNanos = TimeUnit.MINUTES.toNanos(Math.max(1, waitTimeoutMinutes)); long startedAt = System.nanoTime(); + long nextStatusAt = startedAt; while (true) { if (currentBranchJobId != null && jobRepository.existsNewerBranchAnalysisJob( @@ -113,7 +115,11 @@ public GateResult awaitTurn( AnalysisLockType.PR_ANALYSIS.name(), branchName, projectId); } - emitWait(consumer, branchName, sourcePrNumber, waitedNanos); + long now = System.nanoTime(); + if (now >= nextStatusAt) { + emitWait(consumer, branchName, sourcePrNumber, waitedNanos); + nextStatusAt = now + WAIT_STATUS_INTERVAL_NANOS; + } if (!pause()) { throw new AnalysisLockedException( AnalysisLockType.PR_ANALYSIS.name(), branchName, projectId); @@ -158,6 +164,7 @@ public void awaitOlderBranchAnalyses( long timeoutNanos = TimeUnit.MINUTES.toNanos(Math.max(1, waitTimeoutMinutes)); long startedAt = System.nanoTime(); + long nextStatusAt = startedAt; while (jobRepository.existsActiveBranchAnalysisJobBefore( projectId, branchName, currentPrJobId)) { long waitedNanos = System.nanoTime() - startedAt; @@ -169,7 +176,11 @@ public void awaitOlderBranchAnalyses( AnalysisLockType.BRANCH_ANALYSIS.name(), branchName, projectId); } - emitBranchWait(consumer, branchName, waitedNanos); + long now = System.nanoTime(); + if (now >= nextStatusAt) { + emitBranchWait(consumer, branchName, waitedNanos); + nextStatusAt = now + WAIT_STATUS_INTERVAL_NANOS; + } if (!pause()) { throw new AnalysisLockedException( AnalysisLockType.BRANCH_ANALYSIS.name(), branchName, projectId); @@ -212,9 +223,13 @@ private void emitWait( event.put("type", "pr_analysis_wait"); event.put("state", "waiting_for_pr_analysis"); event.put("message", sourcePrNumber == null - ? "Waiting for earlier PR analyses targeting " + branchName + " to finish" - : "Waiting for PR #" + sourcePrNumber + " analysis targeting " - + branchName + " to finish"); + ? "Branch analysis for " + branchName + + " is waiting for earlier PR analyses targeting that branch" + : "Post-merge branch analysis for " + branchName + + " is waiting for PR #" + sourcePrNumber + + " analysis to finish"); + event.put("waitingJobType", JobType.BRANCH_ANALYSIS.name()); + event.put("blockingJobType", JobType.PR_ANALYSIS.name()); event.put("branchName", branchName); event.put("waitedSeconds", TimeUnit.NANOSECONDS.toSeconds(waitedNanos)); if (sourcePrNumber != null) { @@ -237,8 +252,10 @@ private void emitBranchWait( consumer.accept(Map.of( "type", "branch_analysis_wait", "state", "waiting_for_target_branch", - "message", "Waiting for the earlier " + branchName - + " update and its RAG publication to finish", + "message", "PR analysis targeting " + branchName + + " is waiting for the earlier branch update and its RAG publication", + "waitingJobType", JobType.PR_ANALYSIS.name(), + "blockingJobType", JobType.BRANCH_ANALYSIS.name(), "branchName", branchName, "waitedSeconds", TimeUnit.NANOSECONDS.toSeconds(waitedNanos))); } catch (Exception e) { diff --git a/java-ecosystem/libs/analysis-engine/src/test/java/org/rostilos/codecrow/analysisengine/processor/analysis/PullRequestAnalysisProcessorTest.java b/java-ecosystem/libs/analysis-engine/src/test/java/org/rostilos/codecrow/analysisengine/processor/analysis/PullRequestAnalysisProcessorTest.java index 8c3af505..76523a4b 100644 --- a/java-ecosystem/libs/analysis-engine/src/test/java/org/rostilos/codecrow/analysisengine/processor/analysis/PullRequestAnalysisProcessorTest.java +++ b/java-ecosystem/libs/analysis-engine/src/test/java/org/rostilos/codecrow/analysisengine/processor/analysis/PullRequestAnalysisProcessorTest.java @@ -15,7 +15,6 @@ import org.rostilos.codecrow.analysisengine.service.AnalysisLockService; import org.rostilos.codecrow.analysisengine.service.PullRequestService; import org.rostilos.codecrow.commitgraph.service.AnalyzedCommitService; -import org.rostilos.codecrow.analysisengine.service.rag.RagOperationsService; import org.rostilos.codecrow.vcsclient.VcsClientProvider; import org.rostilos.codecrow.analysisengine.service.vcs.VcsAiClientService; import org.rostilos.codecrow.analysisengine.service.vcs.VcsReportingService; @@ -86,9 +85,6 @@ class PullRequestAnalysisProcessorTest { @Mock private AstScopeEnricher astScopeEnricher; - @Mock - private RagOperationsService ragOperationsService; - @Mock private ApplicationEventPublisher eventPublisher; @@ -142,7 +138,6 @@ void setUp() { fileSnapshotService, prIssueTrackingService, astScopeEnricher, - ragOperationsService, eventPublisher); } @@ -757,17 +752,6 @@ void observerFailureDoesNotChangeReviewOutcome() throws Exception { "message", "waiting")); return Optional.of("lock-key-123"); }); - when(ragOperationsService.ensureRagIndexUpToDate( - any(), anyString(), any())) - .thenAnswer(invocation -> { - @SuppressWarnings("unchecked") - java.util.function.Consumer> progress = - invocation.getArgument(2); - progress.accept(Map.of( - "type", "status", - "state", "rag_update")); - return true; - }); when(aiAnalysisClient.performAnalysis(any(), any())).thenAnswer(invocation -> { @SuppressWarnings("unchecked") java.util.function.Consumer> progress = @@ -784,7 +768,6 @@ void observerFailureDoesNotChangeReviewOutcome() throws Exception { verify(reportingService).postAnalysisResults( eq(codeAnalysis), any(), anyLong(), any(), any()); verify(observer).accept(argThat(event -> "lock_wait".equals(event.get("type")))); - verify(observer).accept(argThat(event -> "rag_update".equals(event.get("state")))); verify(observer).accept(argThat(event -> "processing".equals(event.get("state")))); } @@ -1268,7 +1251,6 @@ void shouldWorkWithoutOptionalDependencies() { fileSnapshotService, prIssueTrackingService, null, // astScopeEnricher - null, // ragOperationsService null // eventPublisher ); diff --git a/java-ecosystem/libs/analysis-engine/src/test/java/org/rostilos/codecrow/analysisengine/service/branch/BranchAnalysisGateServiceTest.java b/java-ecosystem/libs/analysis-engine/src/test/java/org/rostilos/codecrow/analysisengine/service/branch/BranchAnalysisGateServiceTest.java index 98bd2709..32ba054b 100644 --- a/java-ecosystem/libs/analysis-engine/src/test/java/org/rostilos/codecrow/analysisengine/service/branch/BranchAnalysisGateServiceTest.java +++ b/java-ecosystem/libs/analysis-engine/src/test/java/org/rostilos/codecrow/analysisengine/service/branch/BranchAnalysisGateServiceTest.java @@ -63,7 +63,7 @@ void branchJobWithoutPrContextWaitsOnlyForEarlierPrJobs() { assertThat(result).isEqualTo(BranchAnalysisGateService.GateResult.READY); verify(jobRepository, times(4)).existsActivePrAnalysisJobBefore(1L, "main", 103L); - verify(consumer, times(3)).accept(org.mockito.ArgumentMatchers.argThat( + verify(consumer).accept(org.mockito.ArgumentMatchers.argThat( event -> "pr_analysis_wait".equals(event.get("type")))); var ordered = inOrder(jobRepository); @@ -111,7 +111,11 @@ void mergeWaitsOnlyForNewestAttemptOfItsOwnPr() { assertThat(result).isEqualTo(BranchAnalysisGateService.GateResult.READY); verify(consumer).accept(org.mockito.ArgumentMatchers.argThat( event -> Long.valueOf(41L).equals(event.get("prNumber")) - && event.get("message").toString().contains("PR #41"))); + && event.get("message").toString().contains("PR #41") + && JobType.BRANCH_ANALYSIS.name().equals( + event.get("waitingJobType")) + && JobType.PR_ANALYSIS.name().equals( + event.get("blockingJobType")))); verify(jobRepository, times(0)).existsActivePrAnalysisJobBefore(1L, "main", 103L); } @@ -164,9 +168,33 @@ void prWaitsOnlyForOlderBranchJobsOnItsTargetBranch() { assertThat(result).isEqualTo(BranchAnalysisGateService.GateResult.READY); verify(jobRepository, times(3)) .existsActiveBranchAnalysisJobBefore(1L, "main", 104L); - verify(consumer, times(2)).accept(org.mockito.ArgumentMatchers.argThat( + verify(consumer).accept(org.mockito.ArgumentMatchers.argThat( event -> "branch_analysis_wait".equals(event.get("type")) - && "main".equals(event.get("branchName")))); + && "main".equals(event.get("branchName")) + && JobType.PR_ANALYSIS.name().equals( + event.get("waitingJobType")) + && JobType.BRANCH_ANALYSIS.name().equals( + event.get("blockingJobType")))); + } + + @Test + void prDoesNotWaitForBranchWorkOnItsSourceOrAnotherTarget() { + Job prJob = job(JobStatus.RUNNING); + ReflectionTestUtils.setField(prJob, "id", 104L); + prJob.setBranchName("main"); + + when(jobRepository.existsActiveBranchAnalysisJobBefore(1L, "main", 104L)) + .thenReturn(false); + + BranchAnalysisGateService.GateResult result = service.awaitDependencies( + 1L, prJob, event -> { }); + + assertThat(result).isEqualTo(BranchAnalysisGateService.GateResult.READY); + verify(jobRepository).existsActiveBranchAnalysisJobBefore(1L, "main", 104L); + verify(jobRepository, times(0)) + .existsActiveBranchAnalysisJobBefore(1L, "1.8.1-rc", 104L); + verify(jobRepository, times(0)) + .existsActiveBranchAnalysisJobBefore(1L, "feature/public-share-links", 104L); } private static Job job(JobStatus status) { diff --git a/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/branch/BranchIndexGenerationBuildService.java b/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/branch/BranchIndexGenerationBuildService.java index 589d1f25..1853e132 100644 --- a/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/branch/BranchIndexGenerationBuildService.java +++ b/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/branch/BranchIndexGenerationBuildService.java @@ -51,7 +51,18 @@ public record PreparedBuild( String collectionTarget, boolean alreadySucceeded, String manifestDigest, - String analysisLockKey) { + String analysisLockKey, + String sourceCollectionTarget) { + + public PreparedBuild( + long operationId, + String collectionTarget, + boolean alreadySucceeded, + String manifestDigest, + String analysisLockKey) { + this(operationId, collectionTarget, alreadySucceeded, manifestDigest, + analysisLockKey, null); + } public PreparedBuild { if (operationId <= 0) { @@ -267,19 +278,36 @@ public Map execute( boolean publishBranchAlias = kind == RagBranchIndexKind.PRIMARY || kind == RagBranchIndexKind.DURABLE; boolean publishLegacyProjectAlias = kind == RagBranchIndexKind.PRIMARY; - Map result = progressEvents == null - ? pipelineClient.indexRepository( - snapshot.toString(), project.getWorkspace().getName(), - project.getNamespace(), branch, revision, includePatterns, - excludePatterns, prepared.collectionTarget(), - false, false) - : pipelineClient.indexRepository( - snapshot.toString(), project.getWorkspace().getName(), - project.getNamespace(), branch, revision, includePatterns, - excludePatterns, prepared.collectionTarget(), - false, false, true, - () -> snapshotOwnershipTransferred.set(true), - progressEvents); + Map result; + if (progressEvents == null) { + result = prepared.sourceCollectionTarget() == null + ? pipelineClient.indexRepository( + snapshot.toString(), project.getWorkspace().getName(), + project.getNamespace(), branch, revision, includePatterns, + excludePatterns, prepared.collectionTarget(), + false, false) + : pipelineClient.indexRepository( + snapshot.toString(), project.getWorkspace().getName(), + project.getNamespace(), branch, revision, includePatterns, + excludePatterns, prepared.collectionTarget(), + false, false, prepared.sourceCollectionTarget()); + } else { + result = prepared.sourceCollectionTarget() == null + ? pipelineClient.indexRepository( + snapshot.toString(), project.getWorkspace().getName(), + project.getNamespace(), branch, revision, includePatterns, + excludePatterns, prepared.collectionTarget(), + false, false, true, + () -> snapshotOwnershipTransferred.set(true), + progressEvents) + : pipelineClient.indexRepository( + snapshot.toString(), project.getWorkspace().getName(), + project.getNamespace(), branch, revision, includePatterns, + excludePatterns, prepared.collectionTarget(), + false, false, true, prepared.sourceCollectionTarget(), + () -> snapshotOwnershipTransferred.set(true), + progressEvents); + } Object manifest = result.get("generation_manifest_sha256"); if (!(manifest instanceof String digest) || digest.isBlank()) { throw new IOException("RAG full branch generation has no manifest digest"); @@ -343,7 +371,8 @@ public static PreparedBuild prepare( registration.generation().getCollectionName(), succeeded, registration.generation().getManifestDigest(), - analysisLockKey); + analysisLockKey, + registration.sourceCollectionTarget()); } private AnalysisLockService.LockLease startAnalysisLockLease(PreparedBuild prepared) diff --git a/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/client/RagPipelineClient.java b/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/client/RagPipelineClient.java index 76b8c3b1..2ee37601 100644 --- a/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/client/RagPipelineClient.java +++ b/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/client/RagPipelineClient.java @@ -132,6 +132,25 @@ public Map indexRepository( String collectionTarget, boolean publishBranchAlias, boolean publishLegacyProjectAlias + ) throws IOException { + return indexRepository( + repoPath, projectWorkspace, projectNamespace, branch, commit, + includePatterns, excludePatterns, collectionTarget, + publishBranchAlias, publishLegacyProjectAlias, (String) null); + } + + public Map indexRepository( + String repoPath, + String projectWorkspace, + String projectNamespace, + String branch, + String commit, + List includePatterns, + List excludePatterns, + String collectionTarget, + boolean publishBranchAlias, + boolean publishLegacyProjectAlias, + String reuseCollectionTarget ) throws IOException { if (!ragEnabled) { log.debug("RAG indexing disabled, skipping repository indexing"); @@ -147,6 +166,9 @@ public Map indexRepository( if (collectionTarget != null && !collectionTarget.isBlank()) { payload.put("collection_target", collectionTarget); } + if (reuseCollectionTarget != null && !reuseCollectionTarget.isBlank()) { + payload.put("reuse_collection_target", reuseCollectionTarget); + } if (publishBranchAlias) { payload.put("publish_branch_alias", true); } @@ -232,6 +254,30 @@ public Map indexRepository( boolean transferRepositoryOwnership, Runnable ownershipAdmissionConsumer, Consumer> progressConsumer + ) throws IOException { + return indexRepository( + repoPath, projectWorkspace, projectNamespace, branch, commit, + includePatterns, excludePatterns, collectionTarget, + publishBranchAlias, publishLegacyProjectAlias, + transferRepositoryOwnership, null, ownershipAdmissionConsumer, + progressConsumer); + } + + public Map indexRepository( + String repoPath, + String projectWorkspace, + String projectNamespace, + String branch, + String commit, + List includePatterns, + List excludePatterns, + String collectionTarget, + boolean publishBranchAlias, + boolean publishLegacyProjectAlias, + boolean transferRepositoryOwnership, + String reuseCollectionTarget, + Runnable ownershipAdmissionConsumer, + Consumer> progressConsumer ) throws IOException { if (!ragEnabled) { log.debug("RAG indexing disabled, skipping repository indexing"); @@ -247,6 +293,9 @@ public Map indexRepository( if (collectionTarget != null && !collectionTarget.isBlank()) { payload.put("collection_target", collectionTarget); } + if (reuseCollectionTarget != null && !reuseCollectionTarget.isBlank()) { + payload.put("reuse_collection_target", reuseCollectionTarget); + } if (publishBranchAlias) { payload.put("publish_branch_alias", true); } diff --git a/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/service/RagBranchIndexRegistryService.java b/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/service/RagBranchIndexRegistryService.java index 8e2dda75..b8e04c5a 100644 --- a/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/service/RagBranchIndexRegistryService.java +++ b/java-ecosystem/libs/rag-engine/src/main/java/org/rostilos/codecrow/ragengine/service/RagBranchIndexRegistryService.java @@ -42,7 +42,16 @@ public record BuildRegistration( RagBranchIndex branchIndex, RagBranchIndexGeneration generation, RagIndexOperation operation, - boolean existingOperation) { + boolean existingOperation, + String sourceCollectionTarget) { + + public BuildRegistration( + RagBranchIndex branchIndex, + RagBranchIndexGeneration generation, + RagIndexOperation operation, + boolean existingOperation) { + this(branchIndex, generation, operation, existingOperation, null); + } } @Transactional @@ -77,7 +86,8 @@ public BuildRegistration registerBuild( branchIndex, operation.getGeneration(), operation, - true); + true, + sourceCollectionTarget(operation.getGeneration())); } RagBranchIndex branchIndex = branchIndexRepository @@ -109,7 +119,18 @@ public BuildRegistration registerBuild( operation.setGeneration(generation); operation = operationRepository.save(operation); - return new BuildRegistration(branchIndex, generation, operation, false); + return new BuildRegistration( + branchIndex, + generation, + operation, + false, + sourceCollectionTarget(generation)); + } + + private static String sourceCollectionTarget( + RagBranchIndexGeneration generation) { + RagBranchIndexGeneration parent = generation.getParentGeneration(); + return parent != null ? parent.getCollectionName() : null; } @Transactional diff --git a/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/branch/BranchIndexBuildAdmissionServiceTest.java b/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/branch/BranchIndexBuildAdmissionServiceTest.java index 75f53dbf..67779ef6 100644 --- a/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/branch/BranchIndexBuildAdmissionServiceTest.java +++ b/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/branch/BranchIndexBuildAdmissionServiceTest.java @@ -55,7 +55,7 @@ void setUp() { operation.setId(30L); operation.setGeneration(generation); registration = new RagBranchIndexRegistryService.BuildRegistration( - branchIndex, generation, operation, false); + branchIndex, generation, operation, false, "source-target"); } @Test @@ -81,6 +81,8 @@ void registersThenAtomicallyLinksAndStartsJobAndOperation() { .isEqualTo("physical-target"); assertThat(admitted.preparedBuild().analysisLockKey()) .isEqualTo("lock-owner-123"); + assertThat(admitted.preparedBuild().sourceCollectionTarget()) + .isEqualTo("source-target"); assertThat(admitted.statusAdmission()).isEqualTo( BranchIndexBuildAdmissionService.ProjectStatusAdmission.UPDATING); InOrder order = inOrder(registryService, jobService, trackingService); diff --git a/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/branch/BranchIndexGenerationBuildServiceTest.java b/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/branch/BranchIndexGenerationBuildServiceTest.java index 555e553f..7c4fad65 100644 --- a/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/branch/BranchIndexGenerationBuildServiceTest.java +++ b/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/branch/BranchIndexGenerationBuildServiceTest.java @@ -71,12 +71,13 @@ void buildsPinnedSnapshotPublishesManifestAndRemovesTemporaryTree() throws Excep project, "develop", RagBranchIndexKind.DURABLE, null, "develop-400", null)) .thenReturn(new RagBranchIndexRegistryService.BuildRegistration( - generation.getBranchIndex(), generation, operation, false)); + generation.getBranchIndex(), generation, operation, false, + "opaque-generation-source")); when(pipelineClient.indexRepository( anyString(), eq("workspace"), eq("namespace"), eq("develop"), eq("develop-400"), eq(List.of("src/**")), eq(List.of("vendor/**")), eq("opaque-generation-target"), - eq(false), eq(false))) + eq(false), eq(false), eq("opaque-generation-source"))) .thenReturn(Map.of( "generation_manifest_sha256", "manifest-400", "document_count", 231, @@ -102,6 +103,11 @@ project, new VcsConnection(), "provider-workspace", "repo", verify(heartbeatService).start(30L); verify(heartbeatScope).close(); verify(registryService).publish(30L, "manifest-400", 231, 400); + verify(pipelineClient).indexRepository( + anyString(), eq("workspace"), eq("namespace"), + eq("develop"), eq("develop-400"), eq(List.of("src/**")), + eq(List.of("vendor/**")), eq("opaque-generation-target"), + eq(false), eq(false), eq("opaque-generation-source")); verify(pipelineClient).publishGenerationAliases( "workspace", "namespace", "develop", "develop-400", "opaque-generation-target", "manifest-400", true, false); @@ -205,11 +211,13 @@ void explicitOperatorRefreshBuildsANewGenerationEvenForTheSameRevision() throws eq(project), eq("develop"), eq(RagBranchIndexKind.DURABLE), isNull(), eq("develop-400"), startsWith("full-snapshot:job:77:"))) .thenReturn(new RagBranchIndexRegistryService.BuildRegistration( - generation.getBranchIndex(), generation, operation, false)); + generation.getBranchIndex(), generation, operation, false, + "opaque-generation-source")); when(pipelineClient.indexRepository( anyString(), anyString(), anyString(), eq("develop"), eq("develop-400"), anyList(), anyList(), eq("opaque-generation-target"), - eq(false), eq(false), eq(true), any(Runnable.class), any())) + eq(false), eq(false), eq(true), eq("opaque-generation-source"), + any(Runnable.class), any())) .thenReturn(Map.of( "generation_manifest_sha256", "fresh-manifest", "document_count", 231, @@ -235,7 +243,7 @@ void explicitOperatorRefreshBuildsANewGenerationEvenForTheSameRevision() throws anyString(), eq("workspace"), eq("namespace"), eq("develop"), eq("develop-400"), anyList(), anyList(), eq("opaque-generation-target"), eq(false), eq(false), eq(true), - any(Runnable.class), any()); + eq("opaque-generation-source"), any(Runnable.class), any()); verify(pipelineClient).publishGenerationAliases( "workspace", "namespace", "develop", "develop-400", "opaque-generation-target", "fresh-manifest", true, false); diff --git a/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/client/RagPipelineClientTest.java b/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/client/RagPipelineClientTest.java index db5537a6..af202787 100644 --- a/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/client/RagPipelineClientTest.java +++ b/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/client/RagPipelineClientTest.java @@ -470,6 +470,25 @@ void testIndexRepository_WithExcludePatterns() throws Exception { assertThat(body).contains("exclude_patterns"); } + @Test + void testIndexRepository_ForwardsPriorGenerationForVectorReuse() throws Exception { + mockWebServer.enqueue(new MockResponse() + .setBody("{\"document_count\":42}") + .addHeader("Content-Type", "application/json")); + + client.indexRepository( + repositoryPath.toString(), "ws", "proj", "main", "abc123", + null, null, "new-generation", false, false, + "prior-generation"); + + RecordedRequest request = mockWebServer.takeRequest(); + Map payload = objectMapper.readValue( + request.getBody().readUtf8(), Map.class); + assertThat(payload) + .containsEntry("collection_target", "new-generation") + .containsEntry("reuse_collection_target", "prior-generation"); + } + @Test void testIndexRepository_WhenDisabled() throws Exception { RagPipelineClient disabledClient = new RagPipelineClient( diff --git a/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/service/RagBranchIndexRegistryServiceTest.java b/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/service/RagBranchIndexRegistryServiceTest.java index 4c1a0534..059e9930 100644 --- a/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/service/RagBranchIndexRegistryServiceTest.java +++ b/java-ecosystem/libs/rag-engine/src/test/java/org/rostilos/codecrow/ragengine/service/RagBranchIndexRegistryServiceTest.java @@ -85,6 +85,7 @@ void registersIdempotentTenantScopedBuildWithoutLeakingBranchName() { .startsWith("cc_w0_p42_b") .doesNotContain("client", "private", "develop"); assertThat(registration.operation().getOperationKey()).hasSize(64); + assertThat(registration.sourceCollectionTarget()).isNull(); } @Test @@ -105,6 +106,8 @@ void publishesNewGenerationAndSupersedesPreviousOneAtomically() { var registration = service.registerBuild( project, "develop", RagBranchIndexKind.DURABLE, "develop-400", "develop-401", "representation"); + assertThat(registration.sourceCollectionTarget()) + .isEqualTo("generation-400"); when(operationRepository.findByIdForUpdate(30L)).thenReturn(Optional.of(registration.operation())); when(branchIndexRepository.findByIdForPublication(10L)) .thenReturn(Optional.of(branchIndex)); diff --git a/python-ecosystem/rag-pipeline/src/rag_pipeline/api/models.py b/python-ecosystem/rag-pipeline/src/rag_pipeline/api/models.py index 7236354a..22028ac7 100644 --- a/python-ecosystem/rag-pipeline/src/rag_pipeline/api/models.py +++ b/python-ecosystem/rag-pipeline/src/rag_pipeline/api/models.py @@ -44,6 +44,7 @@ class IndexRequest(BaseModel): pattern=r"^[0-9a-f]{64}$", ) collection_target: Optional[str] = Field(default=None, min_length=1) + reuse_collection_target: Optional[str] = Field(default=None, min_length=1) publish_branch_alias: bool = False publish_legacy_project_alias: bool = False preserve_other_branches: bool = False diff --git a/python-ecosystem/rag-pipeline/src/rag_pipeline/api/routers/index.py b/python-ecosystem/rag-pipeline/src/rag_pipeline/api/routers/index.py index dd298266..96b325ef 100644 --- a/python-ecosystem/rag-pipeline/src/rag_pipeline/api/routers/index.py +++ b/python-ecosystem/rag-pipeline/src/rag_pipeline/api/routers/index.py @@ -324,6 +324,11 @@ def index_repository(request: IndexRequest, background_tasks: BackgroundTasks): optional_generation_args["source_tree_sha256"] = source_tree_sha256 if isinstance(collection_target, str) and collection_target: optional_generation_args["collection_target"] = collection_target + reuse_collection_target = getattr(request, "reuse_collection_target", None) + if isinstance(reuse_collection_target, str) and reuse_collection_target: + optional_generation_args["reuse_collection_target"] = ( + reuse_collection_target + ) if getattr(request, "publish_branch_alias", False) is True: optional_generation_args["publish_branch_alias"] = True if getattr(request, "publish_legacy_project_alias", False) is True: @@ -400,6 +405,10 @@ def run_index() -> None: optional_generation_args["collection_target"] = ( request.collection_target ) + if request.reuse_collection_target: + optional_generation_args["reuse_collection_target"] = ( + request.reuse_collection_target + ) if getattr(request, "publish_branch_alias", False) is True: optional_generation_args["publish_branch_alias"] = True if getattr(request, "publish_legacy_project_alias", False) is True: diff --git a/python-ecosystem/rag-pipeline/src/rag_pipeline/core/index_manager/indexer.py b/python-ecosystem/rag-pipeline/src/rag_pipeline/core/index_manager/indexer.py index 7bb58dc7..05bcddab 100644 --- a/python-ecosystem/rag-pipeline/src/rag_pipeline/core/index_manager/indexer.py +++ b/python-ecosystem/rag-pipeline/src/rag_pipeline/core/index_manager/indexer.py @@ -432,7 +432,11 @@ def estimate_repository_size( ) is FileDisposition.FULL ] file_count = len(file_list) - logger.info(f"Found {file_count} files for estimation") + logger.info( + "RAG capacity scan found %s repository files " + "(this is not the LLM review scope)", + file_count, + ) if file_count == 0: return 0, 0 @@ -497,6 +501,7 @@ def index_repository( source_tree=None, seal_generation: bool = False, publication_aliases: Optional[List[str]] = None, + reuse_collection_name: Optional[str] = None, operation_id: Optional[str] = None, activation_guard: Optional[Callable[[], None]] = None, progress_callback: Optional[Callable[[dict], None]] = None, @@ -571,6 +576,7 @@ def report_progress( actual_old_collection = None if old_collection_exists: actual_old_collection = self.collection_manager.resolve_alias(alias_name) or alias_name + vector_reuse_collection = reuse_collection_name or actual_old_collection # Get file list repository_file_list = list( @@ -851,7 +857,7 @@ def report_progress( workspace, project, branch, - reuse_collection_name=actual_old_collection, + reuse_collection_name=vector_reuse_collection, operation_id=operation_id, metrics=embedding_metrics, ) @@ -996,7 +1002,7 @@ def report_progress( workspace, project, branch, - reuse_collection_name=actual_old_collection, + reuse_collection_name=vector_reuse_collection, operation_id=operation_id, metrics=embedding_metrics, ) diff --git a/python-ecosystem/rag-pipeline/src/rag_pipeline/core/index_manager/manager.py b/python-ecosystem/rag-pipeline/src/rag_pipeline/core/index_manager/manager.py index 3a909d0f..24b23fc6 100644 --- a/python-ecosystem/rag-pipeline/src/rag_pipeline/core/index_manager/manager.py +++ b/python-ecosystem/rag-pipeline/src/rag_pipeline/core/index_manager/manager.py @@ -333,6 +333,7 @@ def index_repository( exclude_patterns: Optional[List[str]] = None, source_tree_sha256: Optional[str] = None, collection_target: Optional[str] = None, + reuse_collection_target: Optional[str] = None, publish_branch_alias: bool = False, publish_legacy_project_alias: bool = False, progress_callback: Optional[Callable[[dict], None]] = None, @@ -363,6 +364,18 @@ def index_repository( publish_branch_alias, publish_legacy_project_alias, ) + reuse_collection_name = None + if reuse_collection_target: + reuse_collection_name = ( + self._collection_manager.resolve_collection_target( + reuse_collection_target + ) + ) + if reuse_collection_name is None: + logger.info( + "Prior RAG generation is unavailable for vector reuse; " + "embedding the target snapshot normally" + ) with self._mutation_coordinator.acquire( workspace, project, @@ -384,6 +397,7 @@ def index_repository( source_tree=source_tree, seal_generation=collection_target is not None, publication_aliases=publication_aliases, + reuse_collection_name=reuse_collection_name, operation_id=lease.token, activation_guard=lease.assert_owned, progress_callback=progress_callback, diff --git a/python-ecosystem/rag-pipeline/tests/test_index_manager.py b/python-ecosystem/rag-pipeline/tests/test_index_manager.py index 832f3085..9823655b 100644 --- a/python-ecosystem/rag-pipeline/tests/test_index_manager.py +++ b/python-ecosystem/rag-pipeline/tests/test_index_manager.py @@ -581,6 +581,51 @@ def test_readable_alias_publication_requires_immutable_generation_target(self): publish_branch_alias=True, ) + @patch( + "rag_pipeline.core.index_manager.manager.verify_repository_source_tree" + ) + def test_exact_snapshot_forwards_resolved_prior_generation_for_vector_reuse( + self, + verify_source_tree, + ): + from rag_pipeline.core.index_manager.manager import RAGIndexManager + + manager = object.__new__(RAGIndexManager) + manager._collection_manager = MagicMock() + manager._collection_manager.resolve_collection_target.return_value = ( + "prior-generation-physical" + ) + manager._indexer = MagicMock() + manager._indexer.index_repository.return_value = MagicMock() + manager._mutation_coordinator = MagicMock() + lease = MagicMock(token="operation-token") + lease.assert_owned = MagicMock() + manager._mutation_coordinator.acquire.return_value.__enter__.return_value = ( + lease + ) + manager._publication_aliases = MagicMock(return_value=[]) + manager._publication_scope = MagicMock(return_value=None) + source_tree = MagicMock(tree_sha256="f" * 64) + verify_source_tree.return_value = source_tree + + manager.index_repository( + repo_path="/tmp/repository", + workspace="workspace", + project="project", + branch="main", + commit="a" * 40, + source_tree_sha256="f" * 64, + collection_target="new-generation", + reuse_collection_target="prior-generation", + ) + + manager._collection_manager.resolve_collection_target.assert_called_once_with( + "prior-generation" + ) + assert manager._indexer.index_repository.call_args.kwargs[ + "reuse_collection_name" + ] == "prior-generation-physical" + @patch("rag_pipeline.core.index_manager.manager.create_embedding_model") @patch("rag_pipeline.core.index_manager.manager.get_embedding_model_info") @patch("rag_pipeline.core.index_manager.manager.QdrantClient") diff --git a/python-ecosystem/rag-pipeline/tests/test_router_index.py b/python-ecosystem/rag-pipeline/tests/test_router_index.py index 43639611..b65b96f4 100644 --- a/python-ecosystem/rag-pipeline/tests/test_router_index.py +++ b/python-ecosystem/rag-pipeline/tests/test_router_index.py @@ -232,6 +232,7 @@ def index_with_progress(**kwargs): req.exclude_patterns = None req.source_tree_sha256 = None req.collection_target = "target" + req.reuse_collection_target = "prior-target" response = index_repository_stream(req) @@ -246,6 +247,9 @@ async def consume(): assert events[0]["estimatedRemainingMs"] == 1200 assert events[1]["type"] == "complete" assert events[1]["result"]["chunk_count"] == 50 + assert im.index_repository.call_args.kwargs[ + "reuse_collection_target" + ] == "prior-target" def test_stream_progress_is_bounded_and_keeps_latest_event(self): from rag_pipeline.api.routers.index import _coalesce_stream_progress @@ -540,6 +544,7 @@ def test_exact_index_forwards_readable_alias_publication(self, mock_get): repo_path="/tmp/repo", workspace="ws", project="proj", branch="develop", commit="a" * 40, collection_target="exact-develop-target", + reuse_collection_target="prior-develop-target", publish_branch_alias=True, ) @@ -549,6 +554,9 @@ def test_exact_index_forwards_readable_alias_publication(self, mock_get): "exact-develop-target" ) assert im.index_repository.call_args.kwargs["publish_branch_alias"] is True + assert im.index_repository.call_args.kwargs["reuse_collection_target"] == ( + "prior-develop-target" + ) class TestAdvanceGeneration: