Skip to content
Open
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 @@ -290,7 +290,8 @@ private void validateGitHubCollection(
ListenableFuture<CollectionInfo> collectionInfoRequest,
int expectedDimensions,
String denseVectorName,
String sparseVectorName) {
String sparseVectorName)
throws InterruptedException {
try {
CollectionInfo collectionInfo = collectionInfoRequest.get(GRPC_TIMEOUT_SECONDS, TimeUnit.SECONDS);

Expand Down Expand Up @@ -345,9 +346,6 @@ private void validateGitHubCollection(
+ payloadIndex.getKey() + "'");
}
}
} catch (InterruptedException _) {
Thread.currentThread().interrupt();
throw new IllegalStateException("GitHub collection validation interrupted for '" + collectionName + "'");
} catch (ExecutionException executionException) {
if (isTransientGrpcFailure(executionException.getCause())) {
throw new GitHubDiscoveryUnavailableException(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,86 @@ void discoveryCancelsCollectionListingAfterLocalTimeout()
assertEquals("pending", discovery.discoveryHealth().getDetails().get("githubCollectionDiscovery"));
}

@Test
void interruptionDuringValidationStaysPending() throws InterruptedException, ExecutionException, TimeoutException {
QdrantClient qdrantClient = mock(QdrantClient.class);
EmbeddingClient embeddingClient = mock(EmbeddingClient.class);
String activeCollection = GENERATION_PREFIX + "openai-java-chat";
when(qdrantClient.listCollectionsAsync(any(java.time.Duration.class)))
.thenReturn(Futures.immediateFuture(java.util.List.of(activeCollection)));
ListenableFuture<CollectionInfo> collectionInfoRequest = mock();
when(qdrantClient.getCollectionInfoAsync(activeCollection)).thenReturn(collectionInfoRequest);
when(embeddingClient.dimensions()).thenReturn(EMBEDDING_DIMENSIONS);
when(collectionInfoRequest.get(anyLong(), eq(TimeUnit.SECONDS)))
.thenThrow(new InterruptedException("validation interrupted"));
when(collectionInfoRequest.isDone()).thenReturn(false);

QdrantGitHubCollectionDiscovery discovery =
new QdrantGitHubCollectionDiscovery(qdrantClient, embeddingClient, new AppProperties());
discovery.discoverGitHubCollections();

verify(collectionInfoRequest).cancel(true);
assertEquals(java.util.List.of(), discovery.getDiscoveredCollections());
assertEquals(Status.DOWN, discovery.discoveryHealth().getStatus());
assertEquals("pending", discovery.discoveryHealth().getDetails().get("githubCollectionDiscovery"));
}

@Test
void interruptionDuringValidationRecoversOnRetry()
throws InterruptedException, ExecutionException, TimeoutException {
QdrantClient qdrantClient = mock(QdrantClient.class);
EmbeddingClient embeddingClient = mock(EmbeddingClient.class);
String activeCollection = GENERATION_PREFIX + "openai-java-chat";
ListenableFuture<CollectionInfo> interruptedCollectionInfoRequest = mock();
when(qdrantClient.listCollectionsAsync(any(java.time.Duration.class)))
.thenReturn(Futures.immediateFuture(java.util.List.of(activeCollection)));
when(qdrantClient.getCollectionInfoAsync(activeCollection))
.thenReturn(interruptedCollectionInfoRequest)
.thenReturn(Futures.immediateFuture(validCollectionInfo(EMBEDDING_DIMENSIONS)));
when(embeddingClient.dimensions()).thenReturn(EMBEDDING_DIMENSIONS);
when(interruptedCollectionInfoRequest.get(anyLong(), eq(TimeUnit.SECONDS)))
.thenThrow(new InterruptedException("validation interrupted"));
when(interruptedCollectionInfoRequest.isDone()).thenReturn(false);

QdrantGitHubCollectionDiscovery discovery =
new QdrantGitHubCollectionDiscovery(qdrantClient, embeddingClient, new AppProperties());
discovery.discoverGitHubCollections();
assertEquals(Status.DOWN, discovery.discoveryHealth().getStatus());
assertEquals("pending", discovery.discoveryHealth().getDetails().get("githubCollectionDiscovery"));
Thread.interrupted();

discovery.retryPendingDiscovery();
assertEquals(java.util.List.of(activeCollection), discovery.getDiscoveredCollections());
assertEquals(Status.UP, discovery.discoveryHealth().getStatus());
assertEquals("ready", discovery.discoveryHealth().getDetails().get("githubCollectionDiscovery"));
}

@Test
void interruptionDuringValidationRestoresInterruptStatus()
throws InterruptedException, ExecutionException, TimeoutException {
QdrantClient qdrantClient = mock(QdrantClient.class);
EmbeddingClient embeddingClient = mock(EmbeddingClient.class);
String activeCollection = GENERATION_PREFIX + "openai-java-chat";
when(qdrantClient.listCollectionsAsync(any(java.time.Duration.class)))
.thenReturn(Futures.immediateFuture(java.util.List.of(activeCollection)));
ListenableFuture<CollectionInfo> collectionInfoRequest = mock();
when(qdrantClient.getCollectionInfoAsync(activeCollection)).thenReturn(collectionInfoRequest);
when(embeddingClient.dimensions()).thenReturn(EMBEDDING_DIMENSIONS);
when(collectionInfoRequest.get(anyLong(), eq(TimeUnit.SECONDS)))
.thenThrow(new InterruptedException("validation interrupted"));
when(collectionInfoRequest.isDone()).thenReturn(false);

QdrantGitHubCollectionDiscovery discovery =
new QdrantGitHubCollectionDiscovery(qdrantClient, embeddingClient, new AppProperties());
try {
discovery.discoverGitHubCollections();
assertTrue(Thread.currentThread().isInterrupted());
} finally {
Thread.interrupted();
}
assertEquals("pending", discovery.discoveryHealth().getDetails().get("githubCollectionDiscovery"));
}

private static CollectionInfo validCollectionInfo(int denseDimensions) {
VectorParamsMap vectorParams = VectorParamsMap.newBuilder()
.putMap(
Expand Down