Skip to content

Commit 7fd576c

Browse files
authored
Merge branch 'master' into search-cluster-health-reporters
2 parents d3ad2b8 + 6cf7595 commit 7fd576c

21 files changed

Lines changed: 602 additions & 112 deletions

File tree

changelog/unreleased/pr-26674.toml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
type = "f"
2+
message = "Fix Kafka inputs being unable to process LZ4-compressed record batches."
3+
4+
pulls = ["26674"]

data-node/pom.xml

Lines changed: 240 additions & 58 deletions
Large diffs are not rendered by default.

graylog-storage-opensearch3/src/main/java/org/graylog/storage/opensearch3/AdminOpensearchClientProvider.java

Lines changed: 38 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,6 @@
2121
import jakarta.inject.Inject;
2222
import jakarta.inject.Singleton;
2323
import org.graylog.security.certutil.ClientCertSslContextFactory;
24-
import org.graylog.security.certutil.ClientCertSslContextFactoryImpl;
2524
import org.graylog.storage.opensearch3.client.CustomAsyncOpenSearchClient;
2625
import org.graylog.storage.opensearch3.client.CustomOpenSearchClient;
2726
import org.graylog2.configuration.IndexerHosts;
@@ -53,6 +52,10 @@
5352
* returned forever, while the underlying transport is hot-swapped through a
5453
* {@link DynamicTransport} when the cert nears expiry. This makes the returned client safe
5554
* to cache in adapter constructors.
55+
*
56+
* <p>Refresh is lazy and driven by actual use: the cert is only rotated when a request is
57+
* about to be sent and the current cert is close to expiry (see {@link #refreshIfNeeded()},
58+
* wired as the transport's pre-request hook).
5659
*/
5760
@Singleton
5861
public class AdminOpensearchClientProvider {
@@ -70,7 +73,6 @@ public class AdminOpensearchClientProvider {
7073

7174
private volatile OfficialOpensearchClient cachedClient;
7275
private volatile DynamicTransport dynamicTransport;
73-
private volatile ScheduledExecutorService drainScheduler;
7476
private volatile Instant currentCertExpiresAt;
7577

7678
@Inject
@@ -89,48 +91,60 @@ public AdminOpensearchClientProvider(ClientCertSslContextFactory sslContextFacto
8991
/**
9092
* Returns the admin client. The same {@link OfficialOpensearchClient} instance is returned
9193
* across the lifetime of this provider; only the internal transport (and underlying cert)
92-
* is rotated when the cert nears expiry.
94+
* is rotated, lazily, when a request is sent through a client whose cert nears expiry.
9395
*/
9496
@Nonnull
9597
public OfficialOpensearchClient getAdminClient() {
96-
if (cachedClient != null && !needsRefresh(clock.instant())) {
97-
return cachedClient;
98-
}
9998
return initOrRefresh();
10099
}
101100

101+
/**
102+
* Rotates the cert if it is close to expiry. Wired as the {@link DynamicTransport}
103+
* pre-request hook so the cached (and therefore never re-fetched) client picks up a fresh
104+
* cert on its next actual use rather than relying on a background timer.
105+
*/
106+
void refreshIfNeeded() {
107+
if (cachedClient != null && needsRefresh(clock.instant())) {
108+
initOrRefresh();
109+
}
110+
}
111+
102112
private synchronized OfficialOpensearchClient initOrRefresh() {
103113
final Instant now = clock.instant();
104114
if (cachedClient != null && !needsRefresh(now)) {
105115
return cachedClient;
106116
}
107117

118+
if (cachedClient == null) {
119+
this.dynamicTransport = new DynamicTransport(buildTransport(), createDrainScheduler(), this::refreshIfNeeded);
120+
this.cachedClient = new OfficialOpensearchClient(
121+
new CustomOpenSearchClient(dynamicTransport),
122+
new CustomAsyncOpenSearchClient(dynamicTransport),
123+
objectMapper);
124+
this.currentCertExpiresAt = now.plus(CERT_LIFETIME);
125+
LOG.info("Built admin OpenSearch client with a {} min cert lifetime.", CERT_LIFETIME.toMinutes());
126+
} else {
127+
dynamicTransport.swap(buildTransport());
128+
this.currentCertExpiresAt = now.plus(CERT_LIFETIME);
129+
LOG.debug("Rotated admin OpenSearch client certificate.");
130+
}
131+
return cachedClient;
132+
}
133+
134+
/**
135+
* Builds a fresh transport authenticated with a newly minted, short-lived admin cert.
136+
*/
137+
private OpenSearchTransport buildTransport() {
108138
try {
109-
final OpenSearchTransport newTransport = sslContextFactory.buildClientCertSslContext(IndexerAdminCertConstants.ADMIN_CN, CERT_LIFETIME)
139+
return sslContextFactory.buildClientCertSslContext(IndexerAdminCertConstants.ADMIN_CN, CERT_LIFETIME)
110140
.map(TransportConfig::clientCertAuth)
111141
.map(certAuth -> transportProvider.buildTransport(hosts, certAuth))
112-
.orElse(transportProvider.buildTransport(hosts));
113-
114-
if (cachedClient == null) {
115-
this.drainScheduler = createDrainScheduler();
116-
this.dynamicTransport = new DynamicTransport(newTransport, drainScheduler);
117-
this.cachedClient = new OfficialOpensearchClient(
118-
new CustomOpenSearchClient(dynamicTransport),
119-
new CustomAsyncOpenSearchClient(dynamicTransport),
120-
objectMapper);
121-
LOG.info("Built admin OpenSearch client with a {} min cert lifetime.", CERT_LIFETIME.toMinutes());
122-
} else {
123-
dynamicTransport.swap(newTransport);
124-
LOG.debug("Rotated admin OpenSearch client certificate.");
125-
}
126-
this.currentCertExpiresAt = now.plus(CERT_LIFETIME);
142+
.orElseGet(() -> transportProvider.buildTransport(hosts));
127143
} catch (RuntimeException e) {
128144
throw e;
129145
} catch (Exception e) {
130146
throw new IllegalStateException("Failed to build admin OpenSearch client", e);
131147
}
132-
133-
return cachedClient;
134148
}
135149

136150
private boolean needsRefresh(Instant now) {

graylog-storage-opensearch3/src/main/java/org/graylog/storage/opensearch3/DynamicTransport.java

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,18 +39,29 @@ public class DynamicTransport implements OpenSearchTransport {
3939

4040
private final AtomicReference<OpenSearchTransport> current;
4141
private final ScheduledExecutorService scheduler;
42+
private final Runnable beforeRequest;
4243
private final AtomicLong swapGeneration = new AtomicLong(0);
4344
private volatile long drainedGeneration = 0;
4445

4546
public DynamicTransport(OpenSearchTransport initial, ScheduledExecutorService scheduler) {
47+
this(initial, scheduler, () -> {});
48+
}
49+
50+
/**
51+
* @param beforeRequest hook run before every request is dispatched. Lets an owner lazily
52+
* rotate the underlying transport (e.g. to refresh an expiring client
53+
* certificate). Must be cheap and idempotent, as it runs on the request path.
54+
*/
55+
public DynamicTransport(OpenSearchTransport initial, ScheduledExecutorService scheduler, Runnable beforeRequest) {
4656
this.current = new AtomicReference<>(initial);
4757
this.scheduler = scheduler;
58+
this.beforeRequest = beforeRequest;
4859
}
4960

5061
public void swap(OpenSearchTransport newTransport) {
5162
final OpenSearchTransport old = current.getAndSet(newTransport);
5263
final long generation = swapGeneration.incrementAndGet();
53-
LOG.info("OpenSearch transport swapped due to node list update (generation {}). Draining old transport.", generation);
64+
LOG.info("OpenSearch transport swapped (generation {}). Draining old transport.", generation);
5465
try {
5566
scheduler.schedule(() -> {
5667
drainedGeneration = generation;
@@ -80,6 +91,7 @@ public <RequestT, ResponseT, ErrorT> ResponseT performRequest(
8091
RequestT request,
8192
Endpoint<RequestT, ResponseT, ErrorT> endpoint,
8293
TransportOptions options) throws IOException {
94+
beforeRequest.run();
8395
try {
8496
return current.get().performRequest(request, endpoint, options);
8597
} catch (IOException e) {
@@ -92,6 +104,11 @@ public <RequestT, ResponseT, ErrorT> CompletableFuture<ResponseT> performRequest
92104
RequestT request,
93105
Endpoint<RequestT, ResponseT, ErrorT> endpoint,
94106
TransportOptions options) {
107+
try {
108+
beforeRequest.run();
109+
} catch (RuntimeException e) {
110+
return CompletableFuture.failedFuture(e);
111+
}
95112
return current.get().performRequestAsync(request, endpoint, options)
96113
.exceptionally(t -> {
97114
final Throwable cause = unwrapCompletionException(t);

graylog-storage-opensearch3/src/test/java/org/graylog/storage/opensearch3/AdminOpensearchClientProviderTest.java

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -114,6 +114,29 @@ void refreshesTransportWhenCertNearsExpiryButKeepsClientReference() {
114114
verify(transportProvider, times(2)).buildTransport(any(), any());
115115
}
116116

117+
@Test
118+
void refreshIfNeededRotatesTransportOnUseOnlyWhenCertNearsExpiry() {
119+
final Instant start = Instant.parse("2026-01-01T00:00:00Z");
120+
final MutableClock clock = new MutableClock(start);
121+
final AdminOpensearchClientProvider provider = newProvider(clock, initializedSslContextFactory());
122+
123+
final OfficialOpensearchClient client = provider.getAdminClient();
124+
verify(transportProvider, times(1)).buildTransport(any(), any());
125+
126+
// A "use" while the cert is still fresh must not mint a new cert.
127+
provider.refreshIfNeeded();
128+
verify(transportProvider, times(1)).buildTransport(any(), any());
129+
130+
// Once the cert nears expiry, the next "use" rotates the transport in place.
131+
clock.advance(AdminOpensearchClientProvider.CERT_LIFETIME);
132+
provider.refreshIfNeeded();
133+
134+
verify(transportProvider, times(2)).buildTransport(any(), any());
135+
assertThat(provider.getAdminClient())
136+
.as("client reference must remain stable across the on-use rotation")
137+
.isSameAs(client);
138+
}
139+
117140
@Test
118141
void requestsCertWithAdminCommonNameAndConfiguredLifetime() {
119142
final ClientCertSslContextFactory sslContextFactory = (commonName, certificateLifetime) -> {

graylog-storage-opensearch3/src/test/java/org/graylog/storage/opensearch3/DynamicTransportTest.java

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import java.util.concurrent.Executors;
2929
import java.util.concurrent.RejectedExecutionException;
3030
import java.util.concurrent.ScheduledExecutorService;
31+
import java.util.concurrent.atomic.AtomicInteger;
3132

3233
import static org.assertj.core.api.Assertions.assertThat;
3334
import static org.assertj.core.api.Assertions.assertThatCode;
@@ -173,6 +174,42 @@ void asyncPreservesOriginalExceptionWhenNotSwapping() {
173174
.isSameAs(originalException);
174175
}
175176

177+
@Test
178+
@SuppressWarnings("unchecked")
179+
void runsBeforeRequestHookAheadOfEachRequest() throws IOException {
180+
final var delegate = mock(OpenSearchTransport.class);
181+
final var endpoint = mock(Endpoint.class);
182+
final var options = mock(TransportOptions.class);
183+
when(delegate.performRequest(any(), any(), any())).thenReturn("result");
184+
when(delegate.performRequestAsync(any(), any(), any())).thenReturn(CompletableFuture.completedFuture("result"));
185+
186+
final var hookCalls = new AtomicInteger();
187+
final var transport = new DynamicTransport(delegate, scheduler, hookCalls::incrementAndGet);
188+
189+
transport.performRequest("req", endpoint, options);
190+
transport.performRequestAsync("req", endpoint, options);
191+
192+
assertThat(hookCalls).hasValue(2);
193+
}
194+
195+
@Test
196+
@SuppressWarnings("unchecked")
197+
void asyncReturnsFailedFutureWhenBeforeRequestHookThrows() {
198+
final var delegate = mock(OpenSearchTransport.class);
199+
final var endpoint = mock(Endpoint.class);
200+
final var options = mock(TransportOptions.class);
201+
final var hookFailure = new IllegalStateException("cert refresh failed");
202+
203+
final var transport = new DynamicTransport(delegate, scheduler, () -> {
204+
throw hookFailure;
205+
});
206+
207+
final CompletableFuture<?> future = transport.performRequestAsync("req", endpoint, options);
208+
assertThat(future).isCompletedExceptionally();
209+
assertThatThrownBy(future::get).isInstanceOf(ExecutionException.class).cause().isSameAs(hookFailure);
210+
verify(delegate, never()).performRequestAsync(any(), any(), any());
211+
}
212+
176213
@Test
177214
void closesOldTransportImmediatelyWhenSchedulerRejects() throws IOException {
178215
final var oldDelegate = mock(OpenSearchTransport.class);

graylog2-server/pom.xml

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -658,6 +658,11 @@
658658
<groupId>com.github.luben</groupId>
659659
<artifactId>zstd-jni</artifactId>
660660
</dependency>
661+
<!-- Explicitly declare lz4 since kafka-clients depends on it and we exclude it from the declaration. -->
662+
<dependency>
663+
<groupId>at.yawk.lz4</groupId>
664+
<artifactId>lz4-java</artifactId>
665+
</dependency>
661666

662667
<dependency>
663668
<groupId>com.github.zafarkhaja</groupId>

graylog2-server/src/main/java/org/graylog/plugins/views/search/ExplainResults.java

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
*/
1717
package org.graylog.plugins.views.search;
1818

19+
import com.fasterxml.jackson.annotation.JsonInclude;
1920
import org.graylog.plugins.views.search.errors.SearchError;
2021
import org.graylog2.indexer.indexset.MongoIndexSet;
2122
import org.graylog2.indexer.ranges.IndexRange;
@@ -32,7 +33,15 @@ public record SearchResult(Map<String, QueryExplainResult> queries) {
3233
public record QueryExplainResult(Map<String, ExplainResult> searchTypes) {
3334
}
3435

35-
public record ExplainResult(String queryString, Set<IndexRangeResult> searchedIndexRanges) {
36+
public record ExplainResult(String queryString, Set<IndexRangeResult> searchedIndexRanges,
37+
@JsonInclude(JsonInclude.Include.NON_NULL) String effectiveQuery) {
38+
public ExplainResult(String queryString, Set<IndexRangeResult> searchedIndexRanges) {
39+
this(queryString, searchedIndexRanges, null);
40+
}
41+
42+
public ExplainResult withEffectiveQuery(final String effectiveQuery) {
43+
return new ExplainResult(queryString, searchedIndexRanges, effectiveQuery);
44+
}
3645
}
3746

3847
public record IndexRangeResult(String indexName, long begin, long end, boolean isWarmTiered,

graylog2-server/src/main/java/org/graylog/plugins/views/search/engine/QueryEngine.java

Lines changed: 28 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,9 @@
2727
import org.graylog.plugins.views.search.QueryResult;
2828
import org.graylog.plugins.views.search.Search;
2929
import org.graylog.plugins.views.search.SearchJob;
30+
import org.graylog.plugins.views.search.SearchType;
31+
import org.graylog.plugins.views.search.searchfilters.EffectiveQueryComposer;
32+
import org.graylog.plugins.views.search.searchfilters.model.UsedSearchFilter;
3033
import org.graylog.plugins.views.search.elasticsearch.ElasticsearchQueryString;
3134
import org.graylog.plugins.views.search.errors.QueryError;
3235
import org.graylog.plugins.views.search.errors.SearchError;
@@ -38,6 +41,7 @@
3841
import org.slf4j.LoggerFactory;
3942

4043
import java.util.Collection;
44+
import java.util.List;
4145
import java.util.Map;
4246
import java.util.Objects;
4347
import java.util.Set;
@@ -62,17 +66,20 @@ public class QueryEngine {
6266
private final Executor dataLakeJobsQueryPool;
6367
private final ElasticsearchBackendProvider elasticsearchBackendProvider;
6468
private final Map<String, QueryBackend<? extends GeneratedQueryContext>> unversionedBackends;
69+
private final EffectiveQueryComposer effectiveQueryComposer;
6570

6671
@Inject
6772
public QueryEngine(Configuration configuration,
6873
ElasticsearchBackendProvider elasticsearchBackendProvider,
6974
Map<String, QueryBackend<? extends GeneratedQueryContext>> unversionedBackends,
7075
Set<QueryMetadataDecorator> queryMetadataDecorators,
71-
QueryParser queryParser) {
76+
QueryParser queryParser,
77+
EffectiveQueryComposer effectiveQueryComposer) {
7278
this.elasticsearchBackendProvider = elasticsearchBackendProvider;
7379
this.unversionedBackends = unversionedBackends;
7480
this.queryMetadataDecorators = queryMetadataDecorators;
7581
this.queryParser = queryParser;
82+
this.effectiveQueryComposer = effectiveQueryComposer;
7683

7784
this.indexerJobsQueryPool = createThreadPool(
7885
configuration.searchQueryEngineIndexerJobsPoolSize(),
@@ -111,12 +118,31 @@ public ExplainResults explain(SearchJob searchJob, Set<SearchError> validationEr
111118
var backend = getBackendForQuery(q);
112119
final GeneratedQueryContext generatedQueryContext = backend.generate(q, Set.of(), timezone);
113120

114-
return backend.explain(searchJob, q, generatedQueryContext);
121+
return withEffectiveQueries(q, backend.explain(searchJob, q, generatedQueryContext));
115122
}));
116123

117124
return new ExplainResults(searchJob.getSearchId(), new ExplainResults.SearchResult(queries), validationErrors);
118125
}
119126

127+
private ExplainResults.QueryExplainResult withEffectiveQueries(final Query query, final ExplainResults.QueryExplainResult queryExplainResult) {
128+
final Map<String, SearchType> searchTypesById = query.searchTypes().stream()
129+
.collect(Collectors.toMap(SearchType::id, s -> s, (a, b) -> a));
130+
final String baseQuery = query.query().queryString();
131+
132+
final Map<String, ExplainResults.ExplainResult> enriched = queryExplainResult.searchTypes().entrySet().stream()
133+
.collect(Collectors.toMap(Map.Entry::getKey, entry -> {
134+
final List<UsedSearchFilter> filters = new java.util.ArrayList<>(
135+
query.filters() != null ? query.filters() : List.of());
136+
final SearchType searchType = searchTypesById.get(entry.getKey());
137+
if (searchType != null && searchType.filters() != null) {
138+
filters.addAll(searchType.filters());
139+
}
140+
return entry.getValue().withEffectiveQuery(effectiveQueryComposer.compose(baseQuery, filters));
141+
}));
142+
143+
return new ExplainResults.QueryExplainResult(enriched);
144+
}
145+
120146
@WithSpan
121147
public SearchJob execute(SearchJob searchJob, Set<SearchError> validationErrors, DateTimeZone timezone) {
122148
final Set<Query> validQueries = searchJob.getSearch().queries()
Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,32 @@
1+
/*
2+
* Copyright (C) 2020 Graylog, Inc.
3+
*
4+
* This program is free software: you can redistribute it and/or modify
5+
* it under the terms of the Server Side Public License, version 1,
6+
* as published by MongoDB, Inc.
7+
*
8+
* This program is distributed in the hope that it will be useful,
9+
* but WITHOUT ANY WARRANTY; without even the implied warranty of
10+
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
11+
* Server Side Public License for more details.
12+
*
13+
* You should have received a copy of the Server Side Public License
14+
* along with this program. If not, see
15+
* <http://www.mongodb.com/licensing/server-side-public-license>.
16+
*/
17+
package org.graylog.plugins.views.search.rest;
18+
19+
import com.fasterxml.jackson.annotation.JsonCreator;
20+
import com.fasterxml.jackson.annotation.JsonProperty;
21+
import org.graylog.plugins.views.search.searchfilters.model.UsedSearchFilter;
22+
23+
import java.util.List;
24+
25+
public record EffectiveQueryRequest(@JsonProperty("query_string") String queryString,
26+
@JsonProperty("filters") List<UsedSearchFilter> filters) {
27+
@JsonCreator
28+
public EffectiveQueryRequest {
29+
queryString = queryString == null ? "" : queryString;
30+
filters = filters == null ? List.of() : filters;
31+
}
32+
}

0 commit comments

Comments
 (0)