Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -246,9 +242,6 @@ public Map<String, Object> process(
project.getId(), request.getPullRequestId());
}

// Only prepare RAG after the exact snapshot/configuration cache has missed.
ensureRagIndexForTargetBranch(project, request.getTargetBranchName(), consumer);

Map<String, Object> aiResponse = aiAnalysisClient.performAnalysis(aiRequest, event -> {
log.debug("Received event from AI client: type={}", event.get("type"));
emitEvent(consumer, event);
Expand Down Expand Up @@ -944,41 +937,6 @@ static List<VcsCommit> selectPrEvidenceCommits(
return selected;
}

/**
* Ensures RAG index is up-to-date for the PR target branch.
* <p>
* 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
* <p>
* 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
* <p>
* 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.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -90,6 +91,7 @@ public GateResult awaitTurn(
Consumer<Map<String, Object>> consumer) {
long timeoutNanos = TimeUnit.MINUTES.toNanos(Math.max(1, waitTimeoutMinutes));
long startedAt = System.nanoTime();
long nextStatusAt = startedAt;

while (true) {
if (currentBranchJobId != null && jobRepository.existsNewerBranchAnalysisJob(
Expand All @@ -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);
Expand Down Expand Up @@ -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;
Expand All @@ -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);
Expand Down Expand Up @@ -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) {
Expand All @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -86,9 +85,6 @@ class PullRequestAnalysisProcessorTest {
@Mock
private AstScopeEnricher astScopeEnricher;

@Mock
private RagOperationsService ragOperationsService;

@Mock
private ApplicationEventPublisher eventPublisher;

Expand Down Expand Up @@ -142,7 +138,6 @@ void setUp() {
fileSnapshotService,
prIssueTrackingService,
astScopeEnricher,
ragOperationsService,
eventPublisher);
}

Expand Down Expand Up @@ -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<Map<String, Object>> 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<Map<String, Object>> progress =
Expand All @@ -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"))));
}

Expand Down Expand Up @@ -1268,7 +1251,6 @@ void shouldWorkWithoutOptionalDependencies() {
fileSnapshotService,
prIssueTrackingService,
null, // astScopeEnricher
null, // ragOperationsService
null // eventPublisher
);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 MEDIUM | Testing

Wait event assertions expect one invocation

The changed assertions use verify(consumer) without an explicit invocation count, which means exactly one invocation in Mockito. In branchJobWithoutPrContextWaitsOnlyForEarlierPrJobs, the repository is configured to return true three times before false, and the gate emits a wait event on each blocking iteration, so the consumer receives three events. The same mistake appears in prWaitsOnlyForOlderBranchJobsOnItsTargetBranch, where two true responses produce two events. These tests therefore fail against the current gate behavior and no longer validate the expected event count. The current gate implementation emits emitWait/emitBranchWait inside each wait-loop iteration (RAG-5af2ab205aa9f737).

💡 Suggested fix

Restore the expected invocation counts (times(3) and times(2)), or use an explicit atLeastOnce() assertion if the exact polling count is intentionally nondeterministic. Keep the event-schema matcher on the counted verification.

View issue in CodeCrow

event -> "pr_analysis_wait".equals(event.get("type"))));

var ordered = inOrder(jobRepository);
Expand Down Expand Up @@ -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);
}

Expand Down Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -267,19 +278,36 @@ public Map<String, Object> execute(
boolean publishBranchAlias = kind == RagBranchIndexKind.PRIMARY
|| kind == RagBranchIndexKind.DURABLE;
boolean publishLegacyProjectAlias = kind == RagBranchIndexKind.PRIMARY;
Map<String, Object> 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<String, Object> 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");
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading