From 1e400bf7aaaa907200521930ed52808a04fd5d67 Mon Sep 17 00:00:00 2001 From: Shinsuke Sugaya Date: Sat, 3 Oct 2026 10:52:36 +0900 Subject: [PATCH] chore(deps): bump OpenSearch to 3.9.0 and sync the forked tree Bump the version to 3.9.0-SNAPSHOT with the dependency versions OpenSearch 3.9.0 pins (Lucene 10.5.1, Jackson 3.2.2, log4j 2.25.5), and bring the forked OpenSearch classes in org.codelibs.fesen.opensearch up to 3.9.0. Of the 1457 forked classes, 23 differ between 3.8.0 and 3.9.0. The changes the client uses are ported: Version.V_3_9_0, the node concurrency limiter statistics, ClusterAdminClient.deleteTask and the validation and limit changes. HttpDeleteTaskAction implements deleteTask over DELETE /_tasks/{task_id}, and HttpNodesStatsAction parses concurrency_limiters. --- README.md | 2 +- pom.xml | 8 +- .../fesen/client/HttpAbstractClient.java | 23 + .../org/codelibs/fesen/client/HttpClient.java | 8 + .../client/action/HttpDeleteTaskAction.java | 74 +++ .../client/action/HttpNodesStatsAction.java | 83 ++++ .../fesen/opensearch/ExceptionsHelper.java | 2 + .../codelibs/fesen/opensearch/Version.java | 10 +- .../action/ActionConcurrencyLimiterStats.java | 260 ++++++++++ .../admin/cluster/node/stats/NodeStats.java | 28 ++ .../cluster/node/stats/NodesStatsRequest.java | 6 +- .../node/tasks/delete/DeleteTaskAction.java | 33 ++ .../node/tasks/delete/DeleteTaskRequest.java | 81 ++++ .../delete/DeleteTaskRequestBuilder.java | 44 ++ .../validate/query/ValidateQueryResponse.java | 22 +- .../cluster/health/ClusterShardHealth.java | 7 +- .../cluster/metadata/IndexMetadata.java | 107 ++++- .../cluster/metadata/IngestionSource.java | 44 ++ .../routing/IndexShardRoutingTable.java | 3 +- .../fesen/opensearch/common/joda/Joda.java | 448 +++++++++++------- .../common/time/DateFormatters.java | 365 +++++++------- .../core/common/bytes/BytesArray.java | 27 +- .../suggest/completion/FuzzyOptions.java | 31 +- .../suggest/completion/RegexOptions.java | 31 +- .../opensearch/threadpool/ThreadPool.java | 21 +- .../transport/client/ClusterAdminClient.java | 34 ++ .../fesen/client/OpenSearch3ClientTest.java | 39 +- .../fesen/client/action/ActionTestUtils.java | 17 + .../action/HttpDeleteTaskActionTest.java | 69 +++ .../action/HttpNodesStatsActionTest.java | 85 ++++ 30 files changed, 1627 insertions(+), 385 deletions(-) create mode 100644 src/main/java/org/codelibs/fesen/client/action/HttpDeleteTaskAction.java create mode 100644 src/main/java/org/codelibs/fesen/opensearch/action/ActionConcurrencyLimiterStats.java create mode 100644 src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskAction.java create mode 100644 src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskRequest.java create mode 100644 src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskRequestBuilder.java create mode 100644 src/test/java/org/codelibs/fesen/client/action/HttpDeleteTaskActionTest.java diff --git a/README.md b/README.md index 199b3791..733c016b 100644 --- a/README.md +++ b/README.md @@ -226,7 +226,7 @@ The integration test suite runs against the following versions using Testcontain | Engine | Versions tested | |---|---| -| OpenSearch | 1.3.20, 2.19.4, 3.8.0 | +| OpenSearch | 1.3.20, 2.19.4, 3.9.0 | | Elasticsearch | 7.17.29, 8.19.11 | `HttpClient` inspects the cluster's root endpoint (`GET /`) to determine which engine and major version it is talking to, and adjusts request/response handling where the two diverge. diff --git a/pom.xml b/pom.xml index 07ae4516..cc500c00 100644 --- a/pom.xml +++ b/pom.xml @@ -5,7 +5,7 @@ fesen-httpclient jar fesen-httpclient - 3.8.1-SNAPSHOT + 3.9.0-SNAPSHOT HTTP client for OpenSearch https://github.com/codelibs/fesen-httpclient 2012 @@ -38,9 +38,9 @@ UTF-8 21 - 10.5.0 - 3.2.1 - 2.25.4 + 10.5.1 + 3.2.2 + 2.25.5 5.12.2 1.21.4 diff --git a/src/main/java/org/codelibs/fesen/client/HttpAbstractClient.java b/src/main/java/org/codelibs/fesen/client/HttpAbstractClient.java index 805c1957..1d4b5612 100644 --- a/src/main/java/org/codelibs/fesen/client/HttpAbstractClient.java +++ b/src/main/java/org/codelibs/fesen/client/HttpAbstractClient.java @@ -76,6 +76,9 @@ import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksRequest; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksRequestBuilder; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksResponse; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskAction; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskRequest; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskRequestBuilder; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskAction; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskRequest; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskRequestBuilder; @@ -936,6 +939,26 @@ public GetTaskRequestBuilder prepareGetTask(final TaskId taskId) { return new GetTaskRequestBuilder(this, GetTaskAction.INSTANCE).setTaskId(taskId); } + @Override + public ActionFuture deleteTask(final DeleteTaskRequest request) { + return execute(DeleteTaskAction.INSTANCE, request); + } + + @Override + public void deleteTask(final DeleteTaskRequest request, final ActionListener listener) { + execute(DeleteTaskAction.INSTANCE, request, listener); + } + + @Override + public DeleteTaskRequestBuilder prepareDeleteTask(final String taskId) { + return prepareDeleteTask(new TaskId(taskId)); + } + + @Override + public DeleteTaskRequestBuilder prepareDeleteTask(final TaskId taskId) { + return new DeleteTaskRequestBuilder(this, DeleteTaskAction.INSTANCE).setTaskId(taskId); + } + @Override public ActionFuture cancelTasks(final CancelTasksRequest request) { return execute(CancelTasksAction.INSTANCE, request); diff --git a/src/main/java/org/codelibs/fesen/client/HttpClient.java b/src/main/java/org/codelibs/fesen/client/HttpClient.java index 59fcb0be..5f95f7f1 100644 --- a/src/main/java/org/codelibs/fesen/client/HttpClient.java +++ b/src/main/java/org/codelibs/fesen/client/HttpClient.java @@ -68,6 +68,7 @@ import org.codelibs.fesen.client.action.HttpDeleteRepositoryAction; import org.codelibs.fesen.client.action.HttpDeleteSnapshotAction; import org.codelibs.fesen.client.action.HttpDeleteStoredScriptAction; +import org.codelibs.fesen.client.action.HttpDeleteTaskAction; import org.codelibs.fesen.client.action.HttpExplainAction; import org.codelibs.fesen.client.action.HttpFieldCapabilitiesAction; import org.codelibs.fesen.client.action.HttpFlushAction; @@ -167,6 +168,8 @@ import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksAction; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksRequest; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksResponse; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskAction; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskRequest; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskAction; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskRequest; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskResponse; @@ -1057,6 +1060,11 @@ public HttpClient(final Settings settings, final ThreadPool threadPool, final Li new HttpWlmStatsAction(this, WlmStatsAction.INSTANCE).execute((WlmStatsRequest) request, actionListener); }); + actions.put(DeleteTaskAction.INSTANCE, (request, listener) -> { + @SuppressWarnings("unchecked") + final ActionListener actionListener = (ActionListener) listener; + new HttpDeleteTaskAction(this, DeleteTaskAction.INSTANCE).execute((DeleteTaskRequest) request, actionListener); + }); actions.put(GetTaskAction.INSTANCE, (request, listener) -> { @SuppressWarnings("unchecked") final ActionListener actionListener = (ActionListener) listener; diff --git a/src/main/java/org/codelibs/fesen/client/action/HttpDeleteTaskAction.java b/src/main/java/org/codelibs/fesen/client/action/HttpDeleteTaskAction.java new file mode 100644 index 00000000..a34e8d05 --- /dev/null +++ b/src/main/java/org/codelibs/fesen/client/action/HttpDeleteTaskAction.java @@ -0,0 +1,74 @@ +/* + * Copyright 2012-2025 CodeLibs Project and the Others. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, + * either express or implied. See the License for the specific language + * governing permissions and limitations under the License. + */ +package org.codelibs.fesen.client.action; + +import org.codelibs.curl.CurlRequest; +import org.codelibs.fesen.client.HttpClient; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskAction; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskRequest; +import org.codelibs.fesen.opensearch.action.support.clustermanager.AcknowledgedResponse; +import org.codelibs.fesen.opensearch.core.action.ActionListener; +import org.codelibs.fesen.opensearch.core.xcontent.XContentParser; + +/** + * Handles the delete task API over HTTP for OpenSearch, + * removing the stored result of a completed task. + */ +public class HttpDeleteTaskAction extends HttpAction { + + /** The delete task action definition. */ + protected final DeleteTaskAction action; + + /** + * Creates a new HttpDeleteTaskAction. + * + * @param client the HTTP client to send requests with + * @param action the delete task action definition + */ + public HttpDeleteTaskAction(final HttpClient client, final DeleteTaskAction action) { + super(client); + this.action = action; + } + + /** + * Executes the delete task request asynchronously and notifies the listener with the result. + * + * @param request the delete task request + * @param listener the listener to notify with the acknowledged response or a failure + */ + public void execute(final DeleteTaskRequest request, final ActionListener listener) { + getCurlRequest(request).execute(response -> { + try (final XContentParser parser = createParser(response)) { + final AcknowledgedResponse deleteTaskResponse = AcknowledgedResponse.fromXContent(parser); + listener.onResponse(deleteTaskResponse); + } catch (final Exception e) { + listener.onFailure(toOpenSearchException(response, e)); + } + }, e -> unwrapOpenSearchException(listener, e)); + } + + /** + * Builds the HTTP request for the delete task request. + * + * @param request the delete task request + * @return the configured curl request + */ + protected CurlRequest getCurlRequest(final DeleteTaskRequest request) { + // RestDeleteTaskAction + final String taskId = request.getTaskId().getNodeId() + ":" + request.getTaskId().getId(); + return client.getCurlRequest(DELETE, "/_tasks/" + taskId); + } +} diff --git a/src/main/java/org/codelibs/fesen/client/action/HttpNodesStatsAction.java b/src/main/java/org/codelibs/fesen/client/action/HttpNodesStatsAction.java index b75b6c85..0942294c 100644 --- a/src/main/java/org/codelibs/fesen/client/action/HttpNodesStatsAction.java +++ b/src/main/java/org/codelibs/fesen/client/action/HttpNodesStatsAction.java @@ -32,6 +32,7 @@ import org.codelibs.fesen.client.HttpClient; import org.codelibs.fesen.client.io.stream.ByteArrayStreamOutput; import org.codelibs.fesen.opensearch.Version; +import org.codelibs.fesen.opensearch.action.ActionConcurrencyLimiterStats; import org.codelibs.fesen.opensearch.action.admin.cluster.node.stats.NodeStats; import org.codelibs.fesen.opensearch.action.admin.cluster.node.stats.NodesStatsAction; import org.codelibs.fesen.opensearch.action.admin.cluster.node.stats.NodesStatsRequest; @@ -258,6 +259,7 @@ protected NodeStats parseNodeStats(final XContentParser parser, final String nod NodeCacheStats nodeCacheStats = null; RemoteStoreNodeStats remoteStoreNodeStats = null; NativeAllocatorPoolStats nativeAllocatorStats = null; + ActionConcurrencyLimiterStats concurrencyLimiterStats = null; long totalEstimatedNativeBytes = -1L; final Map attributes = new HashMap<>(); XContentParser.Token token; @@ -375,6 +377,8 @@ protected NodeStats parseNodeStats(final XContentParser parser, final String nod consumeObject(parser); } else if ("remote_store".equals(fieldName)) { remoteStoreNodeStats = parseRemoteStoreNodeStats(parser); + } else if ("concurrency_limiters".equals(fieldName)) { + concurrencyLimiterStats = parseConcurrencyLimiterStats(parser); } else { consumeObject(parser); } @@ -438,6 +442,7 @@ protected NodeStats parseNodeStats(final XContentParser parser, final String nod nodeCacheStats, // remoteStoreNodeStats, // nativeAllocatorStats, // + concurrencyLimiterStats, // totalEstimatedNativeBytes); } @@ -1486,6 +1491,84 @@ protected RemoteStoreNodeStats parseRemoteStoreNodeStats(final XContentParser pa } } + /** + * Parses the {@code concurrency_limiters} object, which holds one snapshot per action alias. + * + * @param parser the content parser + * @return the concurrency limiter statistics + * @throws IOException if parsing fails + */ + protected ActionConcurrencyLimiterStats parseConcurrencyLimiterStats(final XContentParser parser) throws IOException { + final List snapshots = new ArrayList<>(); + String alias = null; + XContentParser.Token token; + while ((token = parser.currentToken()) != XContentParser.Token.END_OBJECT) { + if (token == XContentParser.Token.FIELD_NAME) { + alias = parser.currentName(); + } else if (token == XContentParser.Token.START_OBJECT) { + parser.nextToken(); + snapshots.add(parseActionLimiterSnapshot(parser, alias)); + } else { + skipNestedValue(parser); + } + parser.nextToken(); + } + return new ActionConcurrencyLimiterStats(snapshots); + } + + /** + * Parses the snapshot of a single action limiter. The round-trip times are left at {@code -1} + * when the response omits them, which is how the server reports that they are unavailable. + * + * @param parser the content parser + * @param alias the action alias the snapshot is keyed by + * @return the limiter snapshot + * @throws IOException if parsing fails + */ + protected ActionConcurrencyLimiterStats.ActionLimiterSnapshot parseActionLimiterSnapshot(final XContentParser parser, + final String alias) throws IOException { + String actionName = null; + String mode = null; + String algorithm = null; + int currentLimit = 0; + int inFlight = 0; + long totalRejected = 0; + long lastRttMillis = -1L; + long rttNoLoadMillis = -1L; + String fieldName = null; + XContentParser.Token token; + while ((token = parser.currentToken()) != XContentParser.Token.END_OBJECT) { + if (token == XContentParser.Token.FIELD_NAME) { + fieldName = parser.currentName(); + } else if (token == XContentParser.Token.VALUE_STRING) { + if ("action_name".equals(fieldName)) { + actionName = parser.text(); + } else if ("mode".equals(fieldName)) { + mode = parser.text(); + } else if ("algorithm".equals(fieldName)) { + algorithm = parser.text(); + } + } else if (token == XContentParser.Token.VALUE_NUMBER) { + if ("current_limit".equals(fieldName)) { + currentLimit = parser.intValue(); + } else if ("in_flight".equals(fieldName)) { + inFlight = parser.intValue(); + } else if ("total_rejected".equals(fieldName)) { + totalRejected = parser.longValue(); + } else if ("last_rtt_millis".equals(fieldName)) { + lastRttMillis = parser.longValue(); + } else if ("rtt_no_load_millis".equals(fieldName)) { + rttNoLoadMillis = parser.longValue(); + } + } else { + skipNestedValue(parser); + } + parser.nextToken(); + } + return new ActionConcurrencyLimiterStats.ActionLimiterSnapshot(alias, actionName, mode, algorithm, currentLimit, inFlight, + totalRejected, lastRttMillis, rttNoLoadMillis); + } + /** * Parses operation statistics from the response content. * diff --git a/src/main/java/org/codelibs/fesen/opensearch/ExceptionsHelper.java b/src/main/java/org/codelibs/fesen/opensearch/ExceptionsHelper.java index 26dc5ca7..0f47b655 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/ExceptionsHelper.java +++ b/src/main/java/org/codelibs/fesen/opensearch/ExceptionsHelper.java @@ -37,6 +37,7 @@ import org.apache.lucene.index.CorruptIndexException; import org.apache.lucene.index.IndexFormatTooNewException; import org.apache.lucene.index.IndexFormatTooOldException; +import org.apache.lucene.search.IndexSearcher; import org.codelibs.fesen.opensearch.common.CheckedRunnable; import org.codelibs.fesen.opensearch.common.CheckedSupplier; import org.codelibs.fesen.opensearch.common.Nullable; @@ -123,6 +124,7 @@ public static RestStatus status(Throwable t) { case InputCoercionException ignored -> RestStatus.BAD_REQUEST; case JsonParseException ignored -> RestStatus.BAD_REQUEST; case NotXContentException ignored -> RestStatus.BAD_REQUEST; + case IndexSearcher.TooManyClauses ignored -> RestStatus.BAD_REQUEST; case OpenSearchRejectedExecutionException ignored -> RestStatus.TOO_MANY_REQUESTS; case null, default -> RestStatus.INTERNAL_SERVER_ERROR; }; diff --git a/src/main/java/org/codelibs/fesen/opensearch/Version.java b/src/main/java/org/codelibs/fesen/opensearch/Version.java index 682feca1..df1f86c7 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/Version.java +++ b/src/main/java/org/codelibs/fesen/opensearch/Version.java @@ -368,10 +368,18 @@ private static String legacyFriendlyIdToString(int versionId) { * The V_3_8_0 constant. */ public static final Version V_3_8_0 = new Version(3080099, org.apache.lucene.util.Version.LUCENE_10_5_0); + /** + * The V_3_8_1 constant. + */ + public static final Version V_3_8_1 = new Version(3080199, org.apache.lucene.util.Version.LUCENE_10_5_0); + /** + * The V_3_9_0 constant. + */ + public static final Version V_3_9_0 = new Version(3090099, org.apache.lucene.util.Version.LUCENE_10_5_1); /** * The CURRENT constant. */ - public static final Version CURRENT = V_3_8_0; + public static final Version CURRENT = V_3_9_0; /** * The identifier to version. diff --git a/src/main/java/org/codelibs/fesen/opensearch/action/ActionConcurrencyLimiterStats.java b/src/main/java/org/codelibs/fesen/opensearch/action/ActionConcurrencyLimiterStats.java new file mode 100644 index 00000000..ac8002e5 --- /dev/null +++ b/src/main/java/org/codelibs/fesen/opensearch/action/ActionConcurrencyLimiterStats.java @@ -0,0 +1,260 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.codelibs.fesen.opensearch.action; + +import org.codelibs.fesen.opensearch.core.common.io.stream.StreamInput; +import org.codelibs.fesen.opensearch.core.common.io.stream.StreamOutput; +import org.codelibs.fesen.opensearch.core.common.io.stream.Writeable; +import org.codelibs.fesen.opensearch.core.xcontent.ToXContentFragment; +import org.codelibs.fesen.opensearch.core.xcontent.XContentBuilder; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** + * Per-node snapshot of all configured adaptive concurrency limiters, + * exposed via the {@code /_nodes/stats} API under the {@code concurrency_limiters} key. + */ +public class ActionConcurrencyLimiterStats implements Writeable, ToXContentFragment { + + /** + * Snapshot for a single action alias. + */ + public static class ActionLimiterSnapshot implements Writeable, ToXContentFragment { + + private static final long RTT_UNAVAILABLE = -1L; + + private final String alias; + private final String actionName; + private final String mode; + private final String algorithm; + private final int currentLimit; + private final int inFlight; + private final long totalRejected; + private final long lastRttMillis; + private final long rttNoLoadMillis; + + /** + * Creates a snapshot of a single action limiter. + * + * @param alias the action alias + * @param actionName the action name + * @param mode the limiter mode + * @param algorithm the limit algorithm + * @param currentLimit the current concurrency limit + * @param inFlight the number of in-flight requests + * @param totalRejected the total number of rejected requests + * @param lastRttMillis the last round-trip time in milliseconds, or {@code -1} if unavailable + * @param rttNoLoadMillis the no-load round-trip time in milliseconds, or {@code -1} if unavailable + */ + public ActionLimiterSnapshot( + String alias, + String actionName, + String mode, + String algorithm, + int currentLimit, + int inFlight, + long totalRejected, + long lastRttMillis, + long rttNoLoadMillis + ) { + this.alias = alias; + this.actionName = actionName; + this.mode = mode; + this.algorithm = algorithm; + this.currentLimit = currentLimit; + this.inFlight = inFlight; + this.totalRejected = totalRejected; + this.lastRttMillis = lastRttMillis; + this.rttNoLoadMillis = rttNoLoadMillis; + } + + /** + * Creates a snapshot by reading it from the given input. + * + * @param in the input to read from + * @throws IOException if an I/O error occurs + */ + public ActionLimiterSnapshot(StreamInput in) throws IOException { + this.alias = in.readString(); + this.actionName = in.readString(); + this.mode = in.readString(); + this.algorithm = in.readString(); + this.currentLimit = in.readVInt(); + this.inFlight = in.readVInt(); + this.totalRejected = in.readVLong(); + this.lastRttMillis = in.readLong(); + this.rttNoLoadMillis = in.readLong(); + } + + @Override + public void writeTo(StreamOutput out) throws IOException { + out.writeString(alias); + out.writeString(actionName); + out.writeString(mode); + out.writeString(algorithm); + out.writeVInt(currentLimit); + out.writeVInt(inFlight); + out.writeVLong(totalRejected); + out.writeLong(lastRttMillis); + out.writeLong(rttNoLoadMillis); + } + + @Override + public XContentBuilder toXContent(XContentBuilder builder, Params params) throws IOException { + builder.startObject(alias); + builder.field("action_name", actionName); + builder.field("mode", mode); + builder.field("algorithm", algorithm); + builder.field("current_limit", currentLimit); + builder.field("in_flight", inFlight); + builder.field("total_rejected", totalRejected); + if (lastRttMillis != RTT_UNAVAILABLE) builder.field("last_rtt_millis", lastRttMillis); + if (rttNoLoadMillis != RTT_UNAVAILABLE) builder.field("rtt_no_load_millis", rttNoLoadMillis); + builder.endObject(); + return builder; + } + + /** + * Returns the action alias. + * + * @return the action alias + */ + public String getAlias() { + return alias; + } + + /** + * Returns the action name. + * + * @return the action name + */ + public String getActionName() { + return actionName; + } + + /** + * Returns the limiter mode. + * + * @return the limiter mode + */ + public String getMode() { + return mode; + } + + /** + * Returns the limit algorithm. + * + * @return the limit algorithm + */ + public String getAlgorithm() { + return algorithm; + } + + /** + * Returns the current concurrency limit. + * + * @return the current concurrency limit + */ + public int getCurrentLimit() { + return currentLimit; + } + + /** + * Returns the number of in-flight requests. + * + * @return the number of in-flight requests + */ + public int getInFlight() { + return inFlight; + } + + /** + * Returns the total number of rejected requests. + * + * @return the total number of rejected requests + */ + public long getTotalRejected() { + return totalRejected; + } + + /** + * Returns the last round-trip time in milliseconds. + * + * @return the last round-trip time in milliseconds + */ + public long getLastRttMillis() { + return lastRttMillis; + } + + /** + * Returns the no-load round-trip time in milliseconds. + * + * @return the no-load round-trip time in milliseconds + */ + public long getRttNoLoadMillis() { + return rttNoLoadMillis; + } + } + + private final List snapshots; + + /** + * Creates the stats from the given limiter snapshots. + * + * @param snapshots the snapshots, one per configured action alias + */ + public ActionConcurrencyLimiterStats(List snapshots) { + this.snapshots = Collections.unmodifiableList(snapshots); + } + + /** + * Creates the stats by reading them from the given input. + * + * @param in the input to read from + * @throws IOException if an I/O error occurs + */ + public ActionConcurrencyLimiterStats(StreamInput in) throws IOException { + int size = in.readVInt(); + List list = new ArrayList<>(size); + for (int i = 0; i < size; i++) { + list.add(new ActionLimiterSnapshot(in)); + } + this.snapshots = Collections.unmodifiableList(list); + } + + @Override + public void writeTo(StreamOutput out) throws IOException { + out.writeVInt(snapshots.size()); + for (ActionLimiterSnapshot snap : snapshots) { + snap.writeTo(out); + } + } + + @Override + public XContentBuilder toXContent(XContentBuilder builder, Params params) throws IOException { + builder.startObject("concurrency_limiters"); + for (ActionLimiterSnapshot snap : snapshots) { + snap.toXContent(builder, params); + } + builder.endObject(); + return builder; + } + + /** + * Returns the limiter snapshots. + * + * @return the snapshots, one per configured action alias + */ + public List getSnapshots() { + return snapshots; + } +} diff --git a/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/stats/NodeStats.java b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/stats/NodeStats.java index 272396b9..62737bb7 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/stats/NodeStats.java +++ b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/stats/NodeStats.java @@ -33,6 +33,7 @@ package org.codelibs.fesen.opensearch.action.admin.cluster.node.stats; import org.codelibs.fesen.opensearch.Version; +import org.codelibs.fesen.opensearch.action.ActionConcurrencyLimiterStats; import org.codelibs.fesen.opensearch.action.support.nodes.BaseNodeResponse; import org.codelibs.fesen.opensearch.cluster.node.DiscoveryNode; import org.codelibs.fesen.opensearch.cluster.node.DiscoveryNodeRole; @@ -180,6 +181,9 @@ public class NodeStats extends BaseNodeResponse implements ToXContentFragment { @Nullable private NativeAllocatorPoolStats nativeAllocatorStats; + @Nullable + private ActionConcurrencyLimiterStats concurrencyLimiterStats; + /** * Process-level native-memory estimate captured on the data node hosting this {@code NodeStats}. * Computed once in {@link org.codelibs.fesen.opensearch.node.NodeService#stats} via @@ -301,6 +305,11 @@ public NodeStats(StreamInput in) throws IOException { // BWC: V_3_7_0 wrote AnalyticsBackendNativeMemoryStats here; read and discard. in.readOptionalWriteable(AnalyticsBackendNativeMemoryStats::new); } + if (in.getVersion().onOrAfter(Version.V_3_9_0)) { + concurrencyLimiterStats = in.readOptionalWriteable(ActionConcurrencyLimiterStats::new); + } else { + concurrencyLimiterStats = null; + } if (in.getVersion().onOrAfter(Version.V_3_7_0)) { totalEstimatedNativeBytes = in.readLong(); } else { @@ -344,6 +353,7 @@ public NodeStats(StreamInput in) throws IOException { * @param nodeCacheStats the node cache stats * @param remoteStoreNodeStats the remote store node stats * @param nativeAllocatorStats the native allocator stats + * @param concurrencyLimiterStats the concurrency limiter stats * @param totalEstimatedNativeBytes the total estimated native bytes */ public NodeStats( @@ -380,6 +390,7 @@ public NodeStats( @Nullable NodeCacheStats nodeCacheStats, @Nullable RemoteStoreNodeStats remoteStoreNodeStats, @Nullable NativeAllocatorPoolStats nativeAllocatorStats, + @Nullable ActionConcurrencyLimiterStats concurrencyLimiterStats, long totalEstimatedNativeBytes ) { super(node); @@ -415,6 +426,7 @@ public NodeStats( this.nodeCacheStats = nodeCacheStats; this.remoteStoreNodeStats = remoteStoreNodeStats; this.nativeAllocatorStats = nativeAllocatorStats; + this.concurrencyLimiterStats = concurrencyLimiterStats; this.totalEstimatedNativeBytes = totalEstimatedNativeBytes; } @@ -735,6 +747,16 @@ public NativeAllocatorPoolStats getNativeAllocatorStats() { return nativeAllocatorStats; } + /** + * Returns the adaptive concurrency limiter stats, or {@code null} if not available. + * + * @return the concurrency limiter stats + */ + @Nullable + public ActionConcurrencyLimiterStats getConcurrencyLimiterStats() { + return concurrencyLimiterStats; + } + /** * Returns the process-level native-memory estimate captured on this node * (RssAnon - JVM heap committed - JVM non-heap committed), or {@code -1} when the probe @@ -821,6 +843,9 @@ public void writeTo(StreamOutput out) throws IOException { // BWC: V_3_7_0 expects AnalyticsBackendNativeMemoryStats here; write null. out.writeOptionalWriteable(null); } + if (out.getVersion().onOrAfter(Version.V_3_9_0)) { + out.writeOptionalWriteable(concurrencyLimiterStats); + } if (out.getVersion().onOrAfter(Version.V_3_7_0)) { out.writeLong(totalEstimatedNativeBytes); } @@ -944,6 +969,9 @@ public XContentBuilder toXContent(XContentBuilder builder, Params params) throws if (getRemoteStoreNodeStats() != null) { getRemoteStoreNodeStats().toXContent(builder, params); } + if (getConcurrencyLimiterStats() != null) { + getConcurrencyLimiterStats().toXContent(builder, params); + } // total_estimated_bytes ≈ RssAnon - JVM heap committed - JVM non-heap committed. // native_memory: unified view of all native memory pools and jemalloc stats. // NativeAllocatorPoolStats now includes jemalloc allocated/resident + all pools. diff --git a/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/stats/NodesStatsRequest.java b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/stats/NodesStatsRequest.java index d5eefa31..15e5443e 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/stats/NodesStatsRequest.java +++ b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/stats/NodesStatsRequest.java @@ -277,7 +277,11 @@ public enum Metric { /** * The NATIVE_MEMORY value. */ - NATIVE_MEMORY("native_memory"); + NATIVE_MEMORY("native_memory"), + /** + * The CONCURRENCY_LIMITER value. + */ + CONCURRENCY_LIMITER("concurrency_limiter"); private String metricName; diff --git a/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskAction.java b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskAction.java new file mode 100644 index 00000000..8a4d10c7 --- /dev/null +++ b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskAction.java @@ -0,0 +1,33 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete; + +import org.codelibs.fesen.opensearch.action.ActionType; +import org.codelibs.fesen.opensearch.action.support.clustermanager.AcknowledgedResponse; + +/** + * ActionType for deleting a stored completed task result. + * + * @opensearch.internal + */ +public class DeleteTaskAction extends ActionType { + + /** + * The INSTANCE constant. + */ + public static final DeleteTaskAction INSTANCE = new DeleteTaskAction(); + /** + * The NAME constant. + */ + public static final String NAME = "cluster:admin/tasks/delete"; + + private DeleteTaskAction() { + super(NAME, AcknowledgedResponse::new); + } +} diff --git a/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskRequest.java b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskRequest.java new file mode 100644 index 00000000..1dfb4d19 --- /dev/null +++ b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskRequest.java @@ -0,0 +1,81 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete; + +import org.codelibs.fesen.opensearch.action.ActionRequest; +import org.codelibs.fesen.opensearch.action.ActionRequestValidationException; +import org.codelibs.fesen.opensearch.common.annotation.PublicApi; +import org.codelibs.fesen.opensearch.core.common.io.stream.StreamInput; +import org.codelibs.fesen.opensearch.core.common.io.stream.StreamOutput; +import org.codelibs.fesen.opensearch.core.tasks.TaskId; + +import java.io.IOException; + +import static org.codelibs.fesen.opensearch.action.ValidateActions.addValidationError; + +/** + * A request to delete a stored completed task result. + * + * @opensearch.api + */ +@PublicApi(since = "3.8.0") +public class DeleteTaskRequest extends ActionRequest { + private TaskId taskId = TaskId.EMPTY_TASK_ID; + + /** + * Creates a new DeleteTaskRequest. + */ + public DeleteTaskRequest() {} + + /** + * Creates a new DeleteTaskRequest by reading it from the given input. + * + * @param in the input to read from + * @throws IOException if an I/O error occurs + */ + public DeleteTaskRequest(StreamInput in) throws IOException { + super(in); + taskId = TaskId.readFromStream(in); + } + + /** + * Returns the TaskId to delete. + * + * @return the task identifier + */ + public TaskId getTaskId() { + return taskId; + } + + /** + * Set the TaskId to delete. Required. + * + * @param taskId the task identifier + * @return this instance + */ + public DeleteTaskRequest setTaskId(TaskId taskId) { + this.taskId = taskId; + return this; + } + + @Override + public ActionRequestValidationException validate() { + ActionRequestValidationException validationException = null; + if (false == getTaskId().isSet()) { + validationException = addValidationError("task id is required", validationException); + } + return validationException; + } + + @Override + public void writeTo(StreamOutput out) throws IOException { + super.writeTo(out); + taskId.writeTo(out); + } +} diff --git a/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskRequestBuilder.java b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskRequestBuilder.java new file mode 100644 index 00000000..da286319 --- /dev/null +++ b/src/main/java/org/codelibs/fesen/opensearch/action/admin/cluster/node/tasks/delete/DeleteTaskRequestBuilder.java @@ -0,0 +1,44 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * + * The OpenSearch Contributors require contributions made to + * this file be licensed under the Apache-2.0 license or a + * compatible open source license. + */ + +package org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete; + +import org.codelibs.fesen.opensearch.action.ActionRequestBuilder; +import org.codelibs.fesen.opensearch.action.support.clustermanager.AcknowledgedResponse; +import org.codelibs.fesen.opensearch.common.annotation.PublicApi; +import org.codelibs.fesen.opensearch.core.tasks.TaskId; +import org.codelibs.fesen.opensearch.transport.client.OpenSearchClient; + +/** + * Builder for the request to delete a stored completed task result. + * + * @opensearch.api + */ +@PublicApi(since = "3.8.0") +public class DeleteTaskRequestBuilder extends ActionRequestBuilder { + /** + * Creates a new DeleteTaskRequestBuilder. + * + * @param client the client + * @param action the action + */ + public DeleteTaskRequestBuilder(OpenSearchClient client, DeleteTaskAction action) { + super(client, action, new DeleteTaskRequest()); + } + + /** + * Set the TaskId to delete. Required. + * + * @param taskId the task identifier + * @return this instance + */ + public final DeleteTaskRequestBuilder setTaskId(TaskId taskId) { + request.setTaskId(taskId); + return this; + } +} diff --git a/src/main/java/org/codelibs/fesen/opensearch/action/admin/indices/validate/query/ValidateQueryResponse.java b/src/main/java/org/codelibs/fesen/opensearch/action/admin/indices/validate/query/ValidateQueryResponse.java index c17b6971..00da2230 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/action/admin/indices/validate/query/ValidateQueryResponse.java +++ b/src/main/java/org/codelibs/fesen/opensearch/action/admin/indices/validate/query/ValidateQueryResponse.java @@ -93,13 +93,31 @@ public class ValidateQueryResponse extends BroadcastResponse { private final List queryExplanations; - ValidateQueryResponse(StreamInput in) throws IOException { + /** + * Deserialization constructor. Public so plugins that intercept + * {@code ValidateQueryAction} can register response readers. + * + * @param in the stream input + * @throws IOException on deserialization failure + */ + public ValidateQueryResponse(StreamInput in) throws IOException { super(in); valid = in.readBoolean(); queryExplanations = in.readList(QueryExplanation::new); } - ValidateQueryResponse( + /** + * Creates a validate query response. Public so plugins that intercept + * {@code ValidateQueryAction} can construct responses. + * + * @param valid whether the query is valid + * @param queryExplanations per-target explanations, or null for none + * @param totalShards total shards the validation ran on + * @param successfulShards successful shards + * @param failedShards failed shards + * @param shardFailures shard failure details + */ + public ValidateQueryResponse( boolean valid, List queryExplanations, int totalShards, diff --git a/src/main/java/org/codelibs/fesen/opensearch/cluster/health/ClusterShardHealth.java b/src/main/java/org/codelibs/fesen/opensearch/cluster/health/ClusterShardHealth.java index d0bcff16..9ba39565 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/cluster/health/ClusterShardHealth.java +++ b/src/main/java/org/codelibs/fesen/opensearch/cluster/health/ClusterShardHealth.java @@ -404,8 +404,11 @@ public static ClusterHealthStatus getInactivePrimaryHealth(final ShardRouting sh RecoverySource.Type recoveryType = shardRouting.recoverySource().getType(); if (unassignedInfo.getLastAllocationStatus() != AllocationStatus.DECIDERS_NO && unassignedInfo.getNumFailedAllocations() == 0 - && (recoveryType == RecoverySource.Type.EMPTY_STORE - || recoveryType == RecoverySource.Type.LOCAL_SHARDS + && (recoveryType == RecoverySource.Type.EMPTY_STORE || recoveryType == RecoverySource.Type.LOCAL_SHARDS + // REMOTE_STORE is deliberately NOT in this list: a primary hydrating from the remote store cannot + // serve queries until recovery completes, so reporting it YELLOW would overstate availability. It + // stays RED while unassigned/initializing and converges to GREEN without operator action. Revisit + // when warm/searchable-remote shards can serve queries directly off the remote store mid-hydration. || recoveryType == RecoverySource.Type.SNAPSHOT)) { return ClusterHealthStatus.YELLOW; } else { diff --git a/src/main/java/org/codelibs/fesen/opensearch/cluster/metadata/IndexMetadata.java b/src/main/java/org/codelibs/fesen/opensearch/cluster/metadata/IndexMetadata.java index 6d3355d0..5f46cfbf 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/cluster/metadata/IndexMetadata.java +++ b/src/main/java/org/codelibs/fesen/opensearch/cluster/metadata/IndexMetadata.java @@ -515,6 +515,82 @@ public Iterator> settings() { Property.Final ); + /** + * The SETTING_REMOTE_STORE_FENCING_ENABLED constant. + */ + public static final String SETTING_REMOTE_STORE_FENCING_ENABLED = "index.remote_store.fencing.enabled"; + + /** + * Used to enable object-store-backed primary fencing for remote-store-backed indices. When enabled, every + * translog upload on the acknowledgement path performs a compare-and-swap on a per-shard fence blob in the + * translog repository, so a stale primary fails itself instead of acknowledging writes that would be lost on + * failover. Requires the translog repository to support conditional writes. + *

+ * Must not be enabled for an index migrating between document replication and remote store: fencing assumes + * every writer acknowledges through the remote translog upload path, and a docrep-backed writer during + * migration never touches the fence, so it could not be fenced by it. Enable only on indices that are + * remote-store-backed end to end. + */ + public static final Setting INDEX_REMOTE_STORE_FENCING_ENABLED_SETTING = Setting.boolSetting( + SETTING_REMOTE_STORE_FENCING_ENABLED, + false, + Property.IndexScope, + Property.Final + ); + + /** + * The SETTING_REMOTE_STORE_AUTO_RESTORE_ENABLED constant. + */ + public static final String SETTING_REMOTE_STORE_AUTO_RESTORE_ENABLED = "index.remote_store.auto_restore.enabled"; + + /** + * Used to enable automatic restore of a remote-store-backed primary whose last copy was lost with the node that + * held it. When enabled, the allocator converts a primary it has proven unrecoverable from any live node's local + * store (the state that otherwise parks at {@code NO_VALID_SHARD_COPY}, RED) into a remote-store recovery on a + * live node, so the shard goes unassigned → recovering → started with no operator action. + *

+ * Requires {@link #INDEX_REMOTE_STORE_FENCING_ENABLED_SETTING}: the trigger fires on the cluster manager's + * membership view, and the departed primary is typically still alive behind a partition, still able to reach the + * object store. The fence's strictly-greater-term takeover is what makes acknowledging past the restored copy's + * restore point impossible; without it the trigger loses acknowledged writes + * ({@code FenceAutoRestore.tla} refutes the ungated design). + */ + public static final Setting INDEX_REMOTE_STORE_AUTO_RESTORE_ENABLED_SETTING = Setting.boolSetting( + SETTING_REMOTE_STORE_AUTO_RESTORE_ENABLED, + false, + new Setting.Validator<>() { + + @Override + public void validate(final Boolean value) {} + + @Override + public void validate(final Boolean value, final Map, Object> settings) { + if (value) { + final Boolean fencingEnabled = (Boolean) settings.get(INDEX_REMOTE_STORE_FENCING_ENABLED_SETTING); + if (fencingEnabled == null || fencingEnabled == false) { + throw new IllegalArgumentException( + "Setting " + + SETTING_REMOTE_STORE_AUTO_RESTORE_ENABLED + + " can only be enabled when " + + SETTING_REMOTE_STORE_FENCING_ENABLED + + " is enabled: the auto-restore trigger fires on the cluster manager's view of node" + + " membership while the departed primary may still be alive, and only the fence" + + " prevents it from acknowledging writes the restored copy will never see" + ); + } + } + } + + @Override + public Iterator> settings() { + final List> settings = Collections.singletonList(INDEX_REMOTE_STORE_FENCING_ENABLED_SETTING); + return settings.iterator(); + } + }, + Property.IndexScope, + Property.Final + ); + /** * Used to specify if the bulk should use adaptive shard selection to select one shard. */ @@ -1273,6 +1349,32 @@ public Iterator> settings() { }, Property.IndexScope, Property.Final) ); + /** + * Setting for the payload decoder type. + * The built-in decoder is {@code xcontent} (JSON). Plugins may register + * additional decoder names (e.g. {@code avro}). + */ + public static final String SETTING_INGESTION_SOURCE_DECODER_TYPE = "index.ingestion_source.decoder_type"; + /** + * The INGESTION_SOURCE_DECODER_TYPE_SETTING constant. + */ + public static final Setting INGESTION_SOURCE_DECODER_TYPE_SETTING = Setting.simpleString( + SETTING_INGESTION_SOURCE_DECODER_TYPE, + "xcontent", // the built-in decoder's name; the decoder itself is node-side and not carried over + Property.IndexScope, + Property.Final + ); + + /** + * Prefix setting for decoder-specific options. Each decoder factory validates + * supported and required keys. Example: + * {@code index.ingestion_source.decoder_settings.schema_registry_url: https://...} + */ + public static final Setting.AffixSetting INGESTION_SOURCE_DECODER_SETTINGS = Setting.prefixKeySetting( + "index.ingestion_source.decoder_settings.", + key -> new Setting<>(key, "", (value) -> value, Property.IndexScope, Property.Final) + ); + /** * Defines the maximum time to wait for lag to catch up during warmup phase. * A value of -1 means warmup is disabled (the default). A value >= 0 enables warmup with that timeout. @@ -2694,8 +2796,9 @@ public IndexMetadata build() { INDEX_TOTAL_REMOTE_CAPABLE_SHARDS_PER_NODE_SETTING.get(settings); final int indexTotalRemoteCapablePrimaryShardsPerNodeLimit = INDEX_TOTAL_REMOTE_CAPABLE_PRIMARY_SHARDS_PER_NODE_SETTING.get(settings); - final boolean isAppendOnlyIndex = INDEX_APPEND_ONLY_ENABLED_SETTING.get(settings) - || PLUGGABLE_DATAFORMAT_ENABLED_SETTING.get(settings); + final boolean isAppendOnlyIndex = INDEX_APPEND_ONLY_ENABLED_SETTING.exists(settings) + ? INDEX_APPEND_ONLY_ENABLED_SETTING.get(settings) + : PLUGGABLE_DATAFORMAT_ENABLED_SETTING.get(settings); final String uuid = settings.get(SETTING_INDEX_UUID, INDEX_UUID_NA_VALUE); diff --git a/src/main/java/org/codelibs/fesen/opensearch/cluster/metadata/IngestionSource.java b/src/main/java/org/codelibs/fesen/opensearch/cluster/metadata/IngestionSource.java index 77070256..5d9d2def 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/cluster/metadata/IngestionSource.java +++ b/src/main/java/org/codelibs/fesen/opensearch/cluster/metadata/IngestionSource.java @@ -21,6 +21,7 @@ import java.util.Objects; import static org.codelibs.fesen.opensearch.cluster.metadata.IndexMetadata.INGESTION_SOURCE_ALL_ACTIVE_INGESTION_SETTING; +import static org.codelibs.fesen.opensearch.cluster.metadata.IndexMetadata.INGESTION_SOURCE_DECODER_TYPE_SETTING; import static org.codelibs.fesen.opensearch.cluster.metadata.IndexMetadata.INGESTION_SOURCE_INTERNAL_QUEUE_SIZE_SETTING; import static org.codelibs.fesen.opensearch.cluster.metadata.IndexMetadata.INGESTION_SOURCE_MAPPER_TYPE_SETTING; import static org.codelibs.fesen.opensearch.cluster.metadata.IndexMetadata.INGESTION_SOURCE_MAX_POLL_SIZE; @@ -48,6 +49,8 @@ public class IngestionSource { private final TimeValue pointerBasedLagUpdateInterval; private final IngestionMessageMapper.MapperType mapperType; private final Map mapperSettings; + private final String decoderType; + private final Map decoderSettings; private final WarmupConfig warmupConfig; private final SourcePartitionStrategy sourcePartitionStrategy; @@ -64,6 +67,8 @@ private IngestionSource( TimeValue pointerBasedLagUpdateInterval, IngestionMessageMapper.MapperType mapperType, Map mapperSettings, + String decoderType, + Map decoderSettings, WarmupConfig warmupConfig, SourcePartitionStrategy sourcePartitionStrategy ) { @@ -79,6 +84,8 @@ private IngestionSource( this.pointerBasedLagUpdateInterval = pointerBasedLagUpdateInterval; this.mapperType = mapperType; this.mapperSettings = mapperSettings != null ? Collections.unmodifiableMap(mapperSettings) : Collections.emptyMap(); + this.decoderType = decoderType; + this.decoderSettings = decoderSettings != null ? Collections.unmodifiableMap(decoderSettings) : Collections.emptyMap(); this.warmupConfig = warmupConfig; this.sourcePartitionStrategy = sourcePartitionStrategy; } @@ -100,6 +107,8 @@ public boolean equals(Object o) { && Objects.equals(pointerBasedLagUpdateInterval, ingestionSource.pointerBasedLagUpdateInterval) && Objects.equals(mapperType, ingestionSource.mapperType) && Objects.equals(mapperSettings, ingestionSource.mapperSettings) + && Objects.equals(decoderType, ingestionSource.decoderType) + && Objects.equals(decoderSettings, ingestionSource.decoderSettings) && Objects.equals(warmupConfig, ingestionSource.warmupConfig) && Objects.equals(sourcePartitionStrategy, ingestionSource.sourcePartitionStrategy); } @@ -119,6 +128,8 @@ public int hashCode() { pointerBasedLagUpdateInterval, mapperType, mapperSettings, + decoderType, + decoderSettings, warmupConfig, sourcePartitionStrategy ); @@ -155,6 +166,11 @@ public String toString() { + '\'' + ", mapperSettings=" + mapperSettings + + ", decoderType='" + + decoderType + + '\'' + + ", decoderSettings=" + + decoderSettings + ", warmupConfig=" + warmupConfig + ", sourcePartitionStrategy='" @@ -283,6 +299,8 @@ public static class Builder { ); private IngestionMessageMapper.MapperType mapperType = INGESTION_SOURCE_MAPPER_TYPE_SETTING.getDefault(Settings.EMPTY); private Map mapperSettings = new HashMap<>(); + private String decoderType = INGESTION_SOURCE_DECODER_TYPE_SETTING.getDefault(Settings.EMPTY); + private Map decoderSettings = new HashMap<>(); private SourcePartitionStrategy sourcePartitionStrategy = INGESTION_SOURCE_PARTITION_STRATEGY_SETTING.getDefault(Settings.EMPTY); // Warmup configuration private TimeValue warmupTimeout = INGESTION_SOURCE_WARMUP_TIMEOUT_SETTING.getDefault(Settings.EMPTY); @@ -313,6 +331,8 @@ public Builder(IngestionSource ingestionSource) { this.pointerBasedLagUpdateInterval = ingestionSource.pointerBasedLagUpdateInterval; this.mapperType = ingestionSource.mapperType; this.mapperSettings = new HashMap<>(ingestionSource.mapperSettings); + this.decoderType = ingestionSource.decoderType; + this.decoderSettings = new HashMap<>(ingestionSource.decoderSettings); this.sourcePartitionStrategy = ingestionSource.sourcePartitionStrategy; // Copy warmup config WarmupConfig wc = ingestionSource.warmupConfig; @@ -453,6 +473,28 @@ public Builder setMapperSettings(Map mapperSettings) { return this; } + /** + * Sets the payload decoder type. + * + * @param decoderType the decoder type + * @return this instance + */ + public Builder setDecoderType(String decoderType) { + this.decoderType = decoderType; + return this; + } + + /** + * Sets the payload decoder settings. + * + * @param decoderSettings the decoder settings + * @return this instance + */ + public Builder setDecoderSettings(Map decoderSettings) { + this.decoderSettings = decoderSettings; + return this; + } + /** * Sets the source partition strategy. * @@ -518,6 +560,8 @@ public IngestionSource build() { pointerBasedLagUpdateInterval, mapperType, mapperSettings, + decoderType, + decoderSettings, warmupConfig, sourcePartitionStrategy ); diff --git a/src/main/java/org/codelibs/fesen/opensearch/cluster/routing/IndexShardRoutingTable.java b/src/main/java/org/codelibs/fesen/opensearch/cluster/routing/IndexShardRoutingTable.java index ad946224..443ecb65 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/cluster/routing/IndexShardRoutingTable.java +++ b/src/main/java/org/codelibs/fesen/opensearch/cluster/routing/IndexShardRoutingTable.java @@ -469,7 +469,8 @@ public static void writeVerifiableTo(IndexShardRoutingTable indexShard, StreamOu } }); // is primary assigned - out.writeBoolean(indexShard.primaryShard().allocationId() != null); + ShardRouting primary = indexShard.primaryShard(); + out.writeBoolean(primary != null && primary.allocationId() != null); out.writeVInt(indexShard.shards.size() - assignedShardCount.get()); } } diff --git a/src/main/java/org/codelibs/fesen/opensearch/common/joda/Joda.java b/src/main/java/org/codelibs/fesen/opensearch/common/joda/Joda.java index f1ff6b2e..4c621319 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/common/joda/Joda.java +++ b/src/main/java/org/codelibs/fesen/opensearch/common/joda/Joda.java @@ -115,182 +115,278 @@ public static JodaDateFormatter forPattern(String input) { } DateTimeFormatter formatter; - if (FormatNames.BASIC_DATE.matches(input)) { - formatter = ISODateTimeFormat.basicDate(); - } else if (FormatNames.BASIC_DATE_TIME.matches(input)) { - formatter = ISODateTimeFormat.basicDateTime(); - } else if (FormatNames.BASIC_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.basicDateTimeNoMillis(); - } else if (FormatNames.BASIC_ORDINAL_DATE.matches(input)) { - formatter = ISODateTimeFormat.basicOrdinalDate(); - } else if (FormatNames.BASIC_ORDINAL_DATE_TIME.matches(input)) { - formatter = ISODateTimeFormat.basicOrdinalDateTime(); - } else if (FormatNames.BASIC_ORDINAL_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.basicOrdinalDateTimeNoMillis(); - } else if (FormatNames.BASIC_TIME.matches(input)) { - formatter = ISODateTimeFormat.basicTime(); - } else if (FormatNames.BASIC_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.basicTimeNoMillis(); - } else if (FormatNames.BASIC_T_TIME.matches(input)) { - formatter = ISODateTimeFormat.basicTTime(); - } else if (FormatNames.BASIC_T_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.basicTTimeNoMillis(); - } else if (FormatNames.BASIC_WEEK_DATE.matches(input)) { - formatter = ISODateTimeFormat.basicWeekDate(); - } else if (FormatNames.BASIC_WEEK_DATE_TIME.matches(input)) { - formatter = ISODateTimeFormat.basicWeekDateTime(); - } else if (FormatNames.BASIC_WEEK_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.basicWeekDateTimeNoMillis(); - } else if (FormatNames.DATE.matches(input)) { - formatter = ISODateTimeFormat.date(); - } else if (FormatNames.DATE_HOUR.matches(input)) { - formatter = ISODateTimeFormat.dateHour(); - } else if (FormatNames.DATE_HOUR_MINUTE.matches(input)) { - formatter = ISODateTimeFormat.dateHourMinute(); - } else if (FormatNames.DATE_HOUR_MINUTE_SECOND.matches(input)) { - formatter = ISODateTimeFormat.dateHourMinuteSecond(); - } else if (FormatNames.DATE_HOUR_MINUTE_SECOND_FRACTION.matches(input)) { - formatter = ISODateTimeFormat.dateHourMinuteSecondFraction(); - } else if (FormatNames.DATE_HOUR_MINUTE_SECOND_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.dateHourMinuteSecondMillis(); - } else if (FormatNames.DATE_OPTIONAL_TIME.matches(input)) { - // in this case, we have a separate parser and printer since the dataOptionalTimeParser can't print - // this sucks we should use the root local by default and not be dependent on the node - return new JodaDateFormatter( - input, - ISODateTimeFormat.dateOptionalTimeParser().withLocale(Locale.ROOT).withZone(DateTimeZone.UTC).withDefaultYear(1970), - ISODateTimeFormat.dateTime().withLocale(Locale.ROOT).withZone(DateTimeZone.UTC).withDefaultYear(1970) - ); - } else if (FormatNames.DATE_TIME.matches(input)) { - formatter = ISODateTimeFormat.dateTime(); - } else if (FormatNames.DATE_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.dateTimeNoMillis(); - } else if (FormatNames.HOUR.matches(input)) { - formatter = ISODateTimeFormat.hour(); - } else if (FormatNames.HOUR_MINUTE.matches(input)) { - formatter = ISODateTimeFormat.hourMinute(); - } else if (FormatNames.HOUR_MINUTE_SECOND.matches(input)) { - formatter = ISODateTimeFormat.hourMinuteSecond(); - } else if (FormatNames.HOUR_MINUTE_SECOND_FRACTION.matches(input)) { - formatter = ISODateTimeFormat.hourMinuteSecondFraction(); - } else if (FormatNames.HOUR_MINUTE_SECOND_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.hourMinuteSecondMillis(); - } else if (FormatNames.ORDINAL_DATE.matches(input)) { - formatter = ISODateTimeFormat.ordinalDate(); - } else if (FormatNames.ORDINAL_DATE_TIME.matches(input)) { - formatter = ISODateTimeFormat.ordinalDateTime(); - } else if (FormatNames.ORDINAL_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.ordinalDateTimeNoMillis(); - } else if (FormatNames.TIME.matches(input)) { - formatter = ISODateTimeFormat.time(); - } else if (FormatNames.TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.timeNoMillis(); - } else if (FormatNames.T_TIME.matches(input)) { - formatter = ISODateTimeFormat.tTime(); - } else if (FormatNames.T_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.tTimeNoMillis(); - } else if (FormatNames.WEEK_DATE.matches(input)) { - formatter = ISODateTimeFormat.weekDate(); - } else if (FormatNames.WEEK_DATE_TIME.matches(input)) { - formatter = ISODateTimeFormat.weekDateTime(); - } else if (FormatNames.WEEK_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = ISODateTimeFormat.weekDateTimeNoMillis(); - } else if (FormatNames.WEEKYEAR.matches(input)) { - getDeprecationLogger().deprecate( - "week_year_format_name", - "Format name \"week_year\" is deprecated and will be removed in a future version. " + "Use \"weekyear\" format instead" - ); - formatter = ISODateTimeFormat.weekyear(); - } else if (FormatNames.WEEK_YEAR.matches(input)) { - formatter = ISODateTimeFormat.weekyear(); - } else if (FormatNames.WEEK_YEAR_WEEK.matches(input)) { - formatter = ISODateTimeFormat.weekyearWeek(); - } else if (FormatNames.WEEKYEAR_WEEK_DAY.matches(input)) { - formatter = ISODateTimeFormat.weekyearWeekDay(); - } else if (FormatNames.YEAR.matches(input)) { - formatter = ISODateTimeFormat.year(); - } else if (FormatNames.YEAR_MONTH.matches(input)) { - formatter = ISODateTimeFormat.yearMonth(); - } else if (FormatNames.YEAR_MONTH_DAY.matches(input)) { - formatter = ISODateTimeFormat.yearMonthDay(); - } else if (FormatNames.EPOCH_SECOND.matches(input)) { - formatter = new DateTimeFormatterBuilder().append(new EpochTimePrinter(false), new EpochTimeParser(false)).toFormatter(); - } else if (FormatNames.EPOCH_MILLIS.matches(input)) { - formatter = new DateTimeFormatterBuilder().append(new EpochTimePrinter(true), new EpochTimeParser(true)).toFormatter(); - // strict date formats here, must be at least 4 digits for year and two for months and two for day - } else if (FormatNames.STRICT_BASIC_WEEK_DATE.matches(input)) { - formatter = StrictISODateTimeFormat.basicWeekDate(); - } else if (FormatNames.STRICT_BASIC_WEEK_DATE_TIME.matches(input)) { - formatter = StrictISODateTimeFormat.basicWeekDateTime(); - } else if (FormatNames.STRICT_BASIC_WEEK_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = StrictISODateTimeFormat.basicWeekDateTimeNoMillis(); - } else if (FormatNames.STRICT_DATE.matches(input)) { - formatter = StrictISODateTimeFormat.date(); - } else if (FormatNames.STRICT_DATE_HOUR.matches(input)) { - formatter = StrictISODateTimeFormat.dateHour(); - } else if (FormatNames.STRICT_DATE_HOUR_MINUTE.matches(input)) { - formatter = StrictISODateTimeFormat.dateHourMinute(); - } else if (FormatNames.STRICT_DATE_HOUR_MINUTE_SECOND.matches(input)) { - formatter = StrictISODateTimeFormat.dateHourMinuteSecond(); - } else if (FormatNames.STRICT_DATE_HOUR_MINUTE_SECOND_FRACTION.matches(input)) { - formatter = StrictISODateTimeFormat.dateHourMinuteSecondFraction(); - } else if (FormatNames.STRICT_DATE_HOUR_MINUTE_SECOND_MILLIS.matches(input)) { - formatter = StrictISODateTimeFormat.dateHourMinuteSecondMillis(); - } else if (FormatNames.STRICT_DATE_OPTIONAL_TIME.matches(input)) { - // in this case, we have a separate parser and printer since the dataOptionalTimeParser can't print - // this sucks we should use the root local by default and not be dependent on the node - return new JodaDateFormatter( - input, - StrictISODateTimeFormat.dateOptionalTimeParser().withLocale(Locale.ROOT).withZone(DateTimeZone.UTC).withDefaultYear(1970), - StrictISODateTimeFormat.dateTime().withLocale(Locale.ROOT).withZone(DateTimeZone.UTC).withDefaultYear(1970) - ); - } else if (FormatNames.STRICT_DATE_TIME.matches(input)) { - formatter = StrictISODateTimeFormat.dateTime(); - } else if (FormatNames.STRICT_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = StrictISODateTimeFormat.dateTimeNoMillis(); - } else if (FormatNames.STRICT_HOUR.matches(input)) { - formatter = StrictISODateTimeFormat.hour(); - } else if (FormatNames.STRICT_HOUR_MINUTE.matches(input)) { - formatter = StrictISODateTimeFormat.hourMinute(); - } else if (FormatNames.STRICT_HOUR_MINUTE_SECOND.matches(input)) { - formatter = StrictISODateTimeFormat.hourMinuteSecond(); - } else if (FormatNames.STRICT_HOUR_MINUTE_SECOND_FRACTION.matches(input)) { - formatter = StrictISODateTimeFormat.hourMinuteSecondFraction(); - } else if (FormatNames.STRICT_HOUR_MINUTE_SECOND_MILLIS.matches(input)) { - formatter = StrictISODateTimeFormat.hourMinuteSecondMillis(); - } else if (FormatNames.STRICT_ORDINAL_DATE.matches(input)) { - formatter = StrictISODateTimeFormat.ordinalDate(); - } else if (FormatNames.STRICT_ORDINAL_DATE_TIME.matches(input)) { - formatter = StrictISODateTimeFormat.ordinalDateTime(); - } else if (FormatNames.STRICT_ORDINAL_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = StrictISODateTimeFormat.ordinalDateTimeNoMillis(); - } else if (FormatNames.STRICT_TIME.matches(input)) { - formatter = StrictISODateTimeFormat.time(); - } else if (FormatNames.STRICT_TIME_NO_MILLIS.matches(input)) { - formatter = StrictISODateTimeFormat.timeNoMillis(); - } else if (FormatNames.STRICT_T_TIME.matches(input)) { - formatter = StrictISODateTimeFormat.tTime(); - } else if (FormatNames.STRICT_T_TIME_NO_MILLIS.matches(input)) { - formatter = StrictISODateTimeFormat.tTimeNoMillis(); - } else if (FormatNames.STRICT_WEEK_DATE.matches(input)) { - formatter = StrictISODateTimeFormat.weekDate(); - } else if (FormatNames.STRICT_WEEK_DATE_TIME.matches(input)) { - formatter = StrictISODateTimeFormat.weekDateTime(); - } else if (FormatNames.STRICT_WEEK_DATE_TIME_NO_MILLIS.matches(input)) { - formatter = StrictISODateTimeFormat.weekDateTimeNoMillis(); - } else if (FormatNames.STRICT_WEEKYEAR.matches(input)) { - formatter = StrictISODateTimeFormat.weekyear(); - } else if (FormatNames.STRICT_WEEKYEAR_WEEK.matches(input)) { - formatter = StrictISODateTimeFormat.weekyearWeek(); - } else if (FormatNames.STRICT_WEEKYEAR_WEEK_DAY.matches(input)) { - formatter = StrictISODateTimeFormat.weekyearWeekDay(); - } else if (FormatNames.STRICT_YEAR.matches(input)) { - formatter = StrictISODateTimeFormat.year(); - } else if (FormatNames.STRICT_YEAR_MONTH.matches(input)) { - formatter = StrictISODateTimeFormat.yearMonth(); - } else if (FormatNames.STRICT_YEAR_MONTH_DAY.matches(input)) { - formatter = StrictISODateTimeFormat.yearMonthDay(); - } else if (Strings.hasLength(input) && input.contains("||")) { + // formatName is already resolved above via FormatNames.forName(input); switch on the enum directly. + if (formatName != null) { + switch (formatName) { + case BASIC_DATE: + formatter = ISODateTimeFormat.basicDate(); + break; + case BASIC_DATE_TIME: + formatter = ISODateTimeFormat.basicDateTime(); + break; + case BASIC_DATE_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.basicDateTimeNoMillis(); + break; + case BASIC_ORDINAL_DATE: + formatter = ISODateTimeFormat.basicOrdinalDate(); + break; + case BASIC_ORDINAL_DATE_TIME: + formatter = ISODateTimeFormat.basicOrdinalDateTime(); + break; + case BASIC_ORDINAL_DATE_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.basicOrdinalDateTimeNoMillis(); + break; + case BASIC_TIME: + formatter = ISODateTimeFormat.basicTime(); + break; + case BASIC_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.basicTimeNoMillis(); + break; + case BASIC_T_TIME: + formatter = ISODateTimeFormat.basicTTime(); + break; + case BASIC_T_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.basicTTimeNoMillis(); + break; + case BASIC_WEEK_DATE: + formatter = ISODateTimeFormat.basicWeekDate(); + break; + case BASIC_WEEK_DATE_TIME: + formatter = ISODateTimeFormat.basicWeekDateTime(); + break; + case BASIC_WEEK_DATE_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.basicWeekDateTimeNoMillis(); + break; + case DATE: + formatter = ISODateTimeFormat.date(); + break; + case DATE_HOUR: + formatter = ISODateTimeFormat.dateHour(); + break; + case DATE_HOUR_MINUTE: + formatter = ISODateTimeFormat.dateHourMinute(); + break; + case DATE_HOUR_MINUTE_SECOND: + formatter = ISODateTimeFormat.dateHourMinuteSecond(); + break; + case DATE_HOUR_MINUTE_SECOND_FRACTION: + formatter = ISODateTimeFormat.dateHourMinuteSecondFraction(); + break; + case DATE_HOUR_MINUTE_SECOND_MILLIS: + formatter = ISODateTimeFormat.dateHourMinuteSecondMillis(); + break; + case DATE_OPTIONAL_TIME: + // in this case, we have a separate parser and printer since the dataOptionalTimeParser can't print + // this sucks we should use the root local by default and not be dependent on the node + return new JodaDateFormatter( + input, + ISODateTimeFormat.dateOptionalTimeParser().withLocale(Locale.ROOT).withZone(DateTimeZone.UTC).withDefaultYear(1970), + ISODateTimeFormat.dateTime().withLocale(Locale.ROOT).withZone(DateTimeZone.UTC).withDefaultYear(1970) + ); + case DATE_TIME: + formatter = ISODateTimeFormat.dateTime(); + break; + case DATE_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.dateTimeNoMillis(); + break; + case HOUR: + formatter = ISODateTimeFormat.hour(); + break; + case HOUR_MINUTE: + formatter = ISODateTimeFormat.hourMinute(); + break; + case HOUR_MINUTE_SECOND: + formatter = ISODateTimeFormat.hourMinuteSecond(); + break; + case HOUR_MINUTE_SECOND_FRACTION: + formatter = ISODateTimeFormat.hourMinuteSecondFraction(); + break; + case HOUR_MINUTE_SECOND_MILLIS: + formatter = ISODateTimeFormat.hourMinuteSecondMillis(); + break; + case ORDINAL_DATE: + formatter = ISODateTimeFormat.ordinalDate(); + break; + case ORDINAL_DATE_TIME: + formatter = ISODateTimeFormat.ordinalDateTime(); + break; + case ORDINAL_DATE_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.ordinalDateTimeNoMillis(); + break; + case TIME: + formatter = ISODateTimeFormat.time(); + break; + case TIME_NO_MILLIS: + formatter = ISODateTimeFormat.timeNoMillis(); + break; + case T_TIME: + formatter = ISODateTimeFormat.tTime(); + break; + case T_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.tTimeNoMillis(); + break; + case WEEK_DATE: + formatter = ISODateTimeFormat.weekDate(); + break; + case WEEK_DATE_TIME: + formatter = ISODateTimeFormat.weekDateTime(); + break; + case WEEK_DATE_TIME_NO_MILLIS: + formatter = ISODateTimeFormat.weekDateTimeNoMillis(); + break; + case WEEKYEAR: + getDeprecationLogger().deprecate( + "week_year_format_name", + "Format name \"week_year\" is deprecated and will be removed in a future version. Use \"weekyear\" format instead" + ); + formatter = ISODateTimeFormat.weekyear(); + break; + case WEEK_YEAR: + formatter = ISODateTimeFormat.weekyear(); + break; + case WEEK_YEAR_WEEK: + formatter = ISODateTimeFormat.weekyearWeek(); + break; + case WEEKYEAR_WEEK_DAY: + formatter = ISODateTimeFormat.weekyearWeekDay(); + break; + case YEAR: + formatter = ISODateTimeFormat.year(); + break; + case YEAR_MONTH: + formatter = ISODateTimeFormat.yearMonth(); + break; + case YEAR_MONTH_DAY: + formatter = ISODateTimeFormat.yearMonthDay(); + break; + case EPOCH_SECOND: + formatter = new DateTimeFormatterBuilder().append(new EpochTimePrinter(false), new EpochTimeParser(false)) + .toFormatter(); + break; + case EPOCH_MILLIS: + formatter = new DateTimeFormatterBuilder().append(new EpochTimePrinter(true), new EpochTimeParser(true)).toFormatter(); + break; + // strict date formats here, must be at least 4 digits for year and two for months and two for day + case STRICT_BASIC_WEEK_DATE: + formatter = StrictISODateTimeFormat.basicWeekDate(); + break; + case STRICT_BASIC_WEEK_DATE_TIME: + formatter = StrictISODateTimeFormat.basicWeekDateTime(); + break; + case STRICT_BASIC_WEEK_DATE_TIME_NO_MILLIS: + formatter = StrictISODateTimeFormat.basicWeekDateTimeNoMillis(); + break; + case STRICT_DATE: + formatter = StrictISODateTimeFormat.date(); + break; + case STRICT_DATE_HOUR: + formatter = StrictISODateTimeFormat.dateHour(); + break; + case STRICT_DATE_HOUR_MINUTE: + formatter = StrictISODateTimeFormat.dateHourMinute(); + break; + case STRICT_DATE_HOUR_MINUTE_SECOND: + formatter = StrictISODateTimeFormat.dateHourMinuteSecond(); + break; + case STRICT_DATE_HOUR_MINUTE_SECOND_FRACTION: + formatter = StrictISODateTimeFormat.dateHourMinuteSecondFraction(); + break; + case STRICT_DATE_HOUR_MINUTE_SECOND_MILLIS: + formatter = StrictISODateTimeFormat.dateHourMinuteSecondMillis(); + break; + case STRICT_DATE_OPTIONAL_TIME: + // in this case, we have a separate parser and printer since the dataOptionalTimeParser can't print + // this sucks we should use the root local by default and not be dependent on the node + return new JodaDateFormatter( + input, + StrictISODateTimeFormat.dateOptionalTimeParser() + .withLocale(Locale.ROOT) + .withZone(DateTimeZone.UTC) + .withDefaultYear(1970), + StrictISODateTimeFormat.dateTime().withLocale(Locale.ROOT).withZone(DateTimeZone.UTC).withDefaultYear(1970) + ); + case STRICT_DATE_TIME: + formatter = StrictISODateTimeFormat.dateTime(); + break; + case STRICT_DATE_TIME_NO_MILLIS: + formatter = StrictISODateTimeFormat.dateTimeNoMillis(); + break; + case STRICT_HOUR: + formatter = StrictISODateTimeFormat.hour(); + break; + case STRICT_HOUR_MINUTE: + formatter = StrictISODateTimeFormat.hourMinute(); + break; + case STRICT_HOUR_MINUTE_SECOND: + formatter = StrictISODateTimeFormat.hourMinuteSecond(); + break; + case STRICT_HOUR_MINUTE_SECOND_FRACTION: + formatter = StrictISODateTimeFormat.hourMinuteSecondFraction(); + break; + case STRICT_HOUR_MINUTE_SECOND_MILLIS: + formatter = StrictISODateTimeFormat.hourMinuteSecondMillis(); + break; + case STRICT_ORDINAL_DATE: + formatter = StrictISODateTimeFormat.ordinalDate(); + break; + case STRICT_ORDINAL_DATE_TIME: + formatter = StrictISODateTimeFormat.ordinalDateTime(); + break; + case STRICT_ORDINAL_DATE_TIME_NO_MILLIS: + formatter = StrictISODateTimeFormat.ordinalDateTimeNoMillis(); + break; + case STRICT_TIME: + formatter = StrictISODateTimeFormat.time(); + break; + case STRICT_TIME_NO_MILLIS: + formatter = StrictISODateTimeFormat.timeNoMillis(); + break; + case STRICT_T_TIME: + formatter = StrictISODateTimeFormat.tTime(); + break; + case STRICT_T_TIME_NO_MILLIS: + formatter = StrictISODateTimeFormat.tTimeNoMillis(); + break; + case STRICT_WEEK_DATE: + formatter = StrictISODateTimeFormat.weekDate(); + break; + case STRICT_WEEK_DATE_TIME: + formatter = StrictISODateTimeFormat.weekDateTime(); + break; + case STRICT_WEEK_DATE_TIME_NO_MILLIS: + formatter = StrictISODateTimeFormat.weekDateTimeNoMillis(); + break; + case STRICT_WEEKYEAR: + formatter = StrictISODateTimeFormat.weekyear(); + break; + case STRICT_WEEKYEAR_WEEK: + formatter = StrictISODateTimeFormat.weekyearWeek(); + break; + case STRICT_WEEKYEAR_WEEK_DAY: + formatter = StrictISODateTimeFormat.weekyearWeekDay(); + break; + case STRICT_YEAR: + formatter = StrictISODateTimeFormat.year(); + break; + case STRICT_YEAR_MONTH: + formatter = StrictISODateTimeFormat.yearMonth(); + break; + case STRICT_YEAR_MONTH_DAY: + formatter = StrictISODateTimeFormat.yearMonthDay(); + break; + default: + // formatName is non-null but has no Joda formatter (e.g. ISO8601, RFC3339_LENIENT, EPOCH_MICROS); + // fall through to the pattern-based path below. + formatter = null; + break; + } + if (formatter != null) { + formatter = formatter.withLocale(Locale.ROOT).withZone(DateTimeZone.UTC).withDefaultYear(1970); + return new JodaDateFormatter(input, formatter, formatter); + } + } + + if (Strings.hasLength(input) && input.contains("||")) { String[] formats = Strings.delimitedListToStringArray(input, "||"); DateTimeParser[] parsers = new DateTimeParser[formats.length]; diff --git a/src/main/java/org/codelibs/fesen/opensearch/common/time/DateFormatters.java b/src/main/java/org/codelibs/fesen/opensearch/common/time/DateFormatters.java index 6b0f731f..d6861f80 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/common/time/DateFormatters.java +++ b/src/main/java/org/codelibs/fesen/opensearch/common/time/DateFormatters.java @@ -2028,188 +2028,193 @@ static DateFormatter forPattern(String input) { ); } - if (FormatNames.ISO8601.matches(input)) { - return ISO_8601; - } else if (FormatNames.BASIC_DATE.matches(input)) { - return BASIC_DATE; - } else if (FormatNames.BASIC_DATE_TIME.matches(input)) { - return BASIC_DATE_TIME; - } else if (FormatNames.BASIC_DATE_TIME_NO_MILLIS.matches(input)) { - return BASIC_DATE_TIME_NO_MILLIS; - } else if (FormatNames.BASIC_ORDINAL_DATE.matches(input)) { - return BASIC_ORDINAL_DATE; - } else if (FormatNames.BASIC_ORDINAL_DATE_TIME.matches(input)) { - return BASIC_ORDINAL_DATE_TIME; - } else if (FormatNames.BASIC_ORDINAL_DATE_TIME_NO_MILLIS.matches(input)) { - return BASIC_ORDINAL_DATE_TIME_NO_MILLIS; - } else if (FormatNames.BASIC_TIME.matches(input)) { - return BASIC_TIME; - } else if (FormatNames.BASIC_TIME_NO_MILLIS.matches(input)) { - return BASIC_TIME_NO_MILLIS; - } else if (FormatNames.BASIC_T_TIME.matches(input)) { - return BASIC_T_TIME; - } else if (FormatNames.BASIC_T_TIME_NO_MILLIS.matches(input)) { - return BASIC_T_TIME_NO_MILLIS; - } else if (FormatNames.BASIC_WEEK_DATE.matches(input)) { - return BASIC_WEEK_DATE; - } else if (FormatNames.BASIC_WEEK_DATE_TIME.matches(input)) { - return BASIC_WEEK_DATE_TIME; - } else if (FormatNames.BASIC_WEEK_DATE_TIME_NO_MILLIS.matches(input)) { - return BASIC_WEEK_DATE_TIME_NO_MILLIS; - } else if (FormatNames.DATE.matches(input)) { - return DATE; - } else if (FormatNames.DATE_HOUR.matches(input)) { - return DATE_HOUR; - } else if (FormatNames.DATE_HOUR_MINUTE.matches(input)) { - return DATE_HOUR_MINUTE; - } else if (FormatNames.DATE_HOUR_MINUTE_SECOND.matches(input)) { - return DATE_HOUR_MINUTE_SECOND; - } else if (FormatNames.DATE_HOUR_MINUTE_SECOND_FRACTION.matches(input)) { - return DATE_HOUR_MINUTE_SECOND_FRACTION; - } else if (FormatNames.DATE_HOUR_MINUTE_SECOND_MILLIS.matches(input)) { - return DATE_HOUR_MINUTE_SECOND_MILLIS; - } else if (FormatNames.DATE_OPTIONAL_TIME.matches(input)) { - return DATE_OPTIONAL_TIME; - } else if (FormatNames.DATE_TIME.matches(input)) { - return DATE_TIME; - } else if (FormatNames.DATE_TIME_NO_MILLIS.matches(input)) { - return DATE_TIME_NO_MILLIS; - } else if (FormatNames.HOUR.matches(input)) { - return HOUR; - } else if (FormatNames.HOUR_MINUTE.matches(input)) { - return HOUR_MINUTE; - } else if (FormatNames.HOUR_MINUTE_SECOND.matches(input)) { - return HOUR_MINUTE_SECOND; - } else if (FormatNames.HOUR_MINUTE_SECOND_FRACTION.matches(input)) { - return HOUR_MINUTE_SECOND_FRACTION; - } else if (FormatNames.HOUR_MINUTE_SECOND_MILLIS.matches(input)) { - return HOUR_MINUTE_SECOND_MILLIS; - } else if (FormatNames.ORDINAL_DATE.matches(input)) { - return ORDINAL_DATE; - } else if (FormatNames.ORDINAL_DATE_TIME.matches(input)) { - return ORDINAL_DATE_TIME; - } else if (FormatNames.ORDINAL_DATE_TIME_NO_MILLIS.matches(input)) { - return ORDINAL_DATE_TIME_NO_MILLIS; - } else if (FormatNames.TIME.matches(input)) { - return TIME; - } else if (FormatNames.TIME_NO_MILLIS.matches(input)) { - return TIME_NO_MILLIS; - } else if (FormatNames.T_TIME.matches(input)) { - return T_TIME; - } else if (FormatNames.T_TIME_NO_MILLIS.matches(input)) { - return T_TIME_NO_MILLIS; - } else if (FormatNames.WEEK_DATE.matches(input)) { - return WEEK_DATE; - } else if (FormatNames.WEEK_DATE_TIME.matches(input)) { - return WEEK_DATE_TIME; - } else if (FormatNames.WEEK_DATE_TIME_NO_MILLIS.matches(input)) { - return WEEK_DATE_TIME_NO_MILLIS; - } else if (FormatNames.WEEK_YEAR.matches(input)) { - deprecationLogger.getOrCompute() - .deprecate( - "week_year_format_name", - "Format name \"week_year\" is deprecated and will be removed in a future version. " + "Use \"weekyear\" format instead" - ); - return WEEK_YEAR; - } else if (FormatNames.WEEKYEAR.matches(input)) { - return WEEKYEAR; - } else if (FormatNames.WEEK_YEAR_WEEK.matches(input)) { - return WEEKYEAR_WEEK; - } else if (FormatNames.WEEKYEAR_WEEK_DAY.matches(input)) { - return WEEKYEAR_WEEK_DAY; - } else if (FormatNames.YEAR.matches(input)) { - return YEAR; - } else if (FormatNames.YEAR_MONTH.matches(input)) { - return YEAR_MONTH; - } else if (FormatNames.YEAR_MONTH_DAY.matches(input)) { - return YEAR_MONTH_DAY; - } else if (FormatNames.EPOCH_SECOND.matches(input)) { - return EpochTime.SECONDS_FORMATTER; - } else if (FormatNames.EPOCH_MILLIS.matches(input)) { - return EpochTime.MILLIS_FORMATTER; - // strict date formats here, must be at least 4 digits for year and two for months and two for day - } else if (FormatNames.EPOCH_MICROS.matches(input)) { - return EpochTime.MICROS_FORMATTER; - } else if (FormatNames.STRICT_BASIC_WEEK_DATE.matches(input)) { - return STRICT_BASIC_WEEK_DATE; - } else if (FormatNames.STRICT_BASIC_WEEK_DATE_TIME.matches(input)) { - return STRICT_BASIC_WEEK_DATE_TIME; - } else if (FormatNames.STRICT_BASIC_WEEK_DATE_TIME_NO_MILLIS.matches(input)) { - return STRICT_BASIC_WEEK_DATE_TIME_NO_MILLIS; - } else if (FormatNames.STRICT_DATE.matches(input)) { - return STRICT_DATE; - } else if (FormatNames.STRICT_DATE_HOUR.matches(input)) { - return STRICT_DATE_HOUR; - } else if (FormatNames.STRICT_DATE_HOUR_MINUTE.matches(input)) { - return STRICT_DATE_HOUR_MINUTE; - } else if (FormatNames.STRICT_DATE_HOUR_MINUTE_SECOND.matches(input)) { - return STRICT_DATE_HOUR_MINUTE_SECOND; - } else if (FormatNames.STRICT_DATE_HOUR_MINUTE_SECOND_FRACTION.matches(input)) { - return STRICT_DATE_HOUR_MINUTE_SECOND_FRACTION; - } else if (FormatNames.STRICT_DATE_HOUR_MINUTE_SECOND_MILLIS.matches(input)) { - return STRICT_DATE_HOUR_MINUTE_SECOND_MILLIS; - } else if (FormatNames.STRICT_DATE_OPTIONAL_TIME.matches(input)) { - return STRICT_DATE_OPTIONAL_TIME; - } else if (FormatNames.STRICT_DATE_OPTIONAL_TIME_NANOS.matches(input)) { - return STRICT_DATE_OPTIONAL_TIME_NANOS; - } else if (FormatNames.STRICT_DATE_TIME.matches(input)) { - return STRICT_DATE_TIME; - } else if (FormatNames.STRICT_DATE_TIME_NO_MILLIS.matches(input)) { - return STRICT_DATE_TIME_NO_MILLIS; - } else if (FormatNames.STRICT_HOUR.matches(input)) { - return STRICT_HOUR; - } else if (FormatNames.STRICT_HOUR_MINUTE.matches(input)) { - return STRICT_HOUR_MINUTE; - } else if (FormatNames.STRICT_HOUR_MINUTE_SECOND.matches(input)) { - return STRICT_HOUR_MINUTE_SECOND; - } else if (FormatNames.STRICT_HOUR_MINUTE_SECOND_FRACTION.matches(input)) { - return STRICT_HOUR_MINUTE_SECOND_FRACTION; - } else if (FormatNames.STRICT_HOUR_MINUTE_SECOND_MILLIS.matches(input)) { - return STRICT_HOUR_MINUTE_SECOND_MILLIS; - } else if (FormatNames.STRICT_ORDINAL_DATE.matches(input)) { - return STRICT_ORDINAL_DATE; - } else if (FormatNames.STRICT_ORDINAL_DATE_TIME.matches(input)) { - return STRICT_ORDINAL_DATE_TIME; - } else if (FormatNames.STRICT_ORDINAL_DATE_TIME_NO_MILLIS.matches(input)) { - return STRICT_ORDINAL_DATE_TIME_NO_MILLIS; - } else if (FormatNames.STRICT_TIME.matches(input)) { - return STRICT_TIME; - } else if (FormatNames.STRICT_TIME_NO_MILLIS.matches(input)) { - return STRICT_TIME_NO_MILLIS; - } else if (FormatNames.STRICT_T_TIME.matches(input)) { - return STRICT_T_TIME; - } else if (FormatNames.STRICT_T_TIME_NO_MILLIS.matches(input)) { - return STRICT_T_TIME_NO_MILLIS; - } else if (FormatNames.STRICT_WEEK_DATE.matches(input)) { - return STRICT_WEEK_DATE; - } else if (FormatNames.STRICT_WEEK_DATE_TIME.matches(input)) { - return STRICT_WEEK_DATE_TIME; - } else if (FormatNames.STRICT_WEEK_DATE_TIME_NO_MILLIS.matches(input)) { - return STRICT_WEEK_DATE_TIME_NO_MILLIS; - } else if (FormatNames.STRICT_WEEKYEAR.matches(input)) { - return STRICT_WEEKYEAR; - } else if (FormatNames.STRICT_WEEKYEAR_WEEK.matches(input)) { - return STRICT_WEEKYEAR_WEEK; - } else if (FormatNames.STRICT_WEEKYEAR_WEEK_DAY.matches(input)) { - return STRICT_WEEKYEAR_WEEK_DAY; - } else if (FormatNames.STRICT_YEAR.matches(input)) { - return STRICT_YEAR; - } else if (FormatNames.STRICT_YEAR_MONTH.matches(input)) { - return STRICT_YEAR_MONTH; - } else if (FormatNames.STRICT_YEAR_MONTH_DAY.matches(input)) { - return STRICT_YEAR_MONTH_DAY; - } else if (FormatNames.RFC3339_LENIENT.matches(input)) { - return RFC3339_LENIENT_DATE_FORMATTER; - } else { - try { - return new JavaDateFormatter( - input, - new DateTimeFormatterBuilder().appendPattern(input).toFormatter(Locale.ROOT).withResolverStyle(ResolverStyle.STRICT) - ); - } catch (IllegalArgumentException e) { - throw new IllegalArgumentException("Invalid format: [" + input + "]: " + e.getMessage(), e); + if (formatName != null) { + switch (formatName) { + case ISO8601: + return ISO_8601; + case BASIC_DATE: + return BASIC_DATE; + case BASIC_DATE_TIME: + return BASIC_DATE_TIME; + case BASIC_DATE_TIME_NO_MILLIS: + return BASIC_DATE_TIME_NO_MILLIS; + case BASIC_ORDINAL_DATE: + return BASIC_ORDINAL_DATE; + case BASIC_ORDINAL_DATE_TIME: + return BASIC_ORDINAL_DATE_TIME; + case BASIC_ORDINAL_DATE_TIME_NO_MILLIS: + return BASIC_ORDINAL_DATE_TIME_NO_MILLIS; + case BASIC_TIME: + return BASIC_TIME; + case BASIC_TIME_NO_MILLIS: + return BASIC_TIME_NO_MILLIS; + case BASIC_T_TIME: + return BASIC_T_TIME; + case BASIC_T_TIME_NO_MILLIS: + return BASIC_T_TIME_NO_MILLIS; + case BASIC_WEEK_DATE: + return BASIC_WEEK_DATE; + case BASIC_WEEK_DATE_TIME: + return BASIC_WEEK_DATE_TIME; + case BASIC_WEEK_DATE_TIME_NO_MILLIS: + return BASIC_WEEK_DATE_TIME_NO_MILLIS; + case DATE: + return DATE; + case DATE_HOUR: + return DATE_HOUR; + case DATE_HOUR_MINUTE: + return DATE_HOUR_MINUTE; + case DATE_HOUR_MINUTE_SECOND: + return DATE_HOUR_MINUTE_SECOND; + case DATE_HOUR_MINUTE_SECOND_FRACTION: + return DATE_HOUR_MINUTE_SECOND_FRACTION; + case DATE_HOUR_MINUTE_SECOND_MILLIS: + return DATE_HOUR_MINUTE_SECOND_MILLIS; + case DATE_OPTIONAL_TIME: + return DATE_OPTIONAL_TIME; + case DATE_TIME: + return DATE_TIME; + case DATE_TIME_NO_MILLIS: + return DATE_TIME_NO_MILLIS; + case HOUR: + return HOUR; + case HOUR_MINUTE: + return HOUR_MINUTE; + case HOUR_MINUTE_SECOND: + return HOUR_MINUTE_SECOND; + case HOUR_MINUTE_SECOND_FRACTION: + return HOUR_MINUTE_SECOND_FRACTION; + case HOUR_MINUTE_SECOND_MILLIS: + return HOUR_MINUTE_SECOND_MILLIS; + case ORDINAL_DATE: + return ORDINAL_DATE; + case ORDINAL_DATE_TIME: + return ORDINAL_DATE_TIME; + case ORDINAL_DATE_TIME_NO_MILLIS: + return ORDINAL_DATE_TIME_NO_MILLIS; + case TIME: + return TIME; + case TIME_NO_MILLIS: + return TIME_NO_MILLIS; + case T_TIME: + return T_TIME; + case T_TIME_NO_MILLIS: + return T_TIME_NO_MILLIS; + case WEEK_DATE: + return WEEK_DATE; + case WEEK_DATE_TIME: + return WEEK_DATE_TIME; + case WEEK_DATE_TIME_NO_MILLIS: + return WEEK_DATE_TIME_NO_MILLIS; + case WEEK_YEAR: + deprecationLogger.getOrCompute() + .deprecate( + "week_year_format_name", + "Format name \"week_year\" is deprecated and will be removed in a future version. Use \"weekyear\" format instead" + ); + return WEEK_YEAR; + case WEEKYEAR: + return WEEKYEAR; + case WEEK_YEAR_WEEK: + return WEEKYEAR_WEEK; + case WEEKYEAR_WEEK_DAY: + return WEEKYEAR_WEEK_DAY; + case YEAR: + return YEAR; + case YEAR_MONTH: + return YEAR_MONTH; + case YEAR_MONTH_DAY: + return YEAR_MONTH_DAY; + case EPOCH_SECOND: + return EpochTime.SECONDS_FORMATTER; + case EPOCH_MILLIS: + return EpochTime.MILLIS_FORMATTER; + // strict date formats here, must be at least 4 digits for year and two for months and two for day + case EPOCH_MICROS: + return EpochTime.MICROS_FORMATTER; + case STRICT_BASIC_WEEK_DATE: + return STRICT_BASIC_WEEK_DATE; + case STRICT_BASIC_WEEK_DATE_TIME: + return STRICT_BASIC_WEEK_DATE_TIME; + case STRICT_BASIC_WEEK_DATE_TIME_NO_MILLIS: + return STRICT_BASIC_WEEK_DATE_TIME_NO_MILLIS; + case STRICT_DATE: + return STRICT_DATE; + case STRICT_DATE_HOUR: + return STRICT_DATE_HOUR; + case STRICT_DATE_HOUR_MINUTE: + return STRICT_DATE_HOUR_MINUTE; + case STRICT_DATE_HOUR_MINUTE_SECOND: + return STRICT_DATE_HOUR_MINUTE_SECOND; + case STRICT_DATE_HOUR_MINUTE_SECOND_FRACTION: + return STRICT_DATE_HOUR_MINUTE_SECOND_FRACTION; + case STRICT_DATE_HOUR_MINUTE_SECOND_MILLIS: + return STRICT_DATE_HOUR_MINUTE_SECOND_MILLIS; + case STRICT_DATE_OPTIONAL_TIME: + return STRICT_DATE_OPTIONAL_TIME; + case STRICT_DATE_OPTIONAL_TIME_NANOS: + return STRICT_DATE_OPTIONAL_TIME_NANOS; + case STRICT_DATE_TIME: + return STRICT_DATE_TIME; + case STRICT_DATE_TIME_NO_MILLIS: + return STRICT_DATE_TIME_NO_MILLIS; + case STRICT_HOUR: + return STRICT_HOUR; + case STRICT_HOUR_MINUTE: + return STRICT_HOUR_MINUTE; + case STRICT_HOUR_MINUTE_SECOND: + return STRICT_HOUR_MINUTE_SECOND; + case STRICT_HOUR_MINUTE_SECOND_FRACTION: + return STRICT_HOUR_MINUTE_SECOND_FRACTION; + case STRICT_HOUR_MINUTE_SECOND_MILLIS: + return STRICT_HOUR_MINUTE_SECOND_MILLIS; + case STRICT_ORDINAL_DATE: + return STRICT_ORDINAL_DATE; + case STRICT_ORDINAL_DATE_TIME: + return STRICT_ORDINAL_DATE_TIME; + case STRICT_ORDINAL_DATE_TIME_NO_MILLIS: + return STRICT_ORDINAL_DATE_TIME_NO_MILLIS; + case STRICT_TIME: + return STRICT_TIME; + case STRICT_TIME_NO_MILLIS: + return STRICT_TIME_NO_MILLIS; + case STRICT_T_TIME: + return STRICT_T_TIME; + case STRICT_T_TIME_NO_MILLIS: + return STRICT_T_TIME_NO_MILLIS; + case STRICT_WEEK_DATE: + return STRICT_WEEK_DATE; + case STRICT_WEEK_DATE_TIME: + return STRICT_WEEK_DATE_TIME; + case STRICT_WEEK_DATE_TIME_NO_MILLIS: + return STRICT_WEEK_DATE_TIME_NO_MILLIS; + case STRICT_WEEKYEAR: + return STRICT_WEEKYEAR; + case STRICT_WEEKYEAR_WEEK: + return STRICT_WEEKYEAR_WEEK; + case STRICT_WEEKYEAR_WEEK_DAY: + return STRICT_WEEKYEAR_WEEK_DAY; + case STRICT_YEAR: + return STRICT_YEAR; + case STRICT_YEAR_MONTH: + return STRICT_YEAR_MONTH; + case STRICT_YEAR_MONTH_DAY: + return STRICT_YEAR_MONTH_DAY; + case RFC3339_LENIENT: + return RFC3339_LENIENT_DATE_FORMATTER; + default: + break; } } + + try { + return new JavaDateFormatter( + input, + new DateTimeFormatterBuilder().appendPattern(input).toFormatter(Locale.ROOT).withResolverStyle(ResolverStyle.STRICT) + ); + } catch (IllegalArgumentException e) { + throw new IllegalArgumentException("Invalid format: [" + input + "]: " + e.getMessage(), e); + } } private static final LocalDate LOCALDATE_EPOCH = LocalDate.of(1970, 1, 1); diff --git a/src/main/java/org/codelibs/fesen/opensearch/core/common/bytes/BytesArray.java b/src/main/java/org/codelibs/fesen/opensearch/core/common/bytes/BytesArray.java index 134064b3..da6e81b7 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/core/common/bytes/BytesArray.java +++ b/src/main/java/org/codelibs/fesen/opensearch/core/common/bytes/BytesArray.java @@ -34,6 +34,7 @@ import org.apache.lucene.util.BitUtil; import org.apache.lucene.util.BytesRef; +import org.apache.lucene.util.UnicodeUtil; import org.codelibs.fesen.opensearch.core.common.io.stream.StreamInput; import java.io.IOException; @@ -51,6 +52,7 @@ public final class BytesArray extends AbstractBytesReference { * The EMPTY constant. */ public static final BytesArray EMPTY = new BytesArray(BytesRef.EMPTY_BYTES, 0, 0); + private static final int MAX_UTF16_LENGTH_FOR_UTF8 = Integer.MAX_VALUE / UnicodeUtil.MAX_UTF8_BYTES_PER_CHAR; private final byte[] bytes; private final int offset; private final int length; @@ -61,7 +63,30 @@ public final class BytesArray extends AbstractBytesReference { * @param bytes the bytes */ public BytesArray(String bytes) { - this(new BytesRef(bytes)); + this(toBytesRef(bytes)); + } + + private static BytesRef toBytesRef(String bytes) { + ensureUTF16LengthIsValidForUTF8Encoding(bytes.length()); + return new BytesRef(bytes); + } + + /** + * Validates that a UTF-16 string can be UTF-8 encoded into a Lucene {@code BytesRef} without integer overflow in + * {@code UnicodeUtil#maxUTF8Length(int)}. Package-private: the only callers are this class and its tests. + * + * @param utf16Length UTF-16 length of the string to encode + */ + static void ensureUTF16LengthIsValidForUTF8Encoding(int utf16Length) { + if ((long) utf16Length * UnicodeUtil.MAX_UTF8_BYTES_PER_CHAR > Integer.MAX_VALUE) { + throw new IllegalArgumentException( + "UTF16 string length [" + + utf16Length + + "] exceeds maximum [" + + MAX_UTF16_LENGTH_FOR_UTF8 + + "] that can be UTF-8 encoded without integer overflow" + ); + } } /** diff --git a/src/main/java/org/codelibs/fesen/opensearch/search/suggest/completion/FuzzyOptions.java b/src/main/java/org/codelibs/fesen/opensearch/search/suggest/completion/FuzzyOptions.java index f9dc27d8..1bac1e12 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/search/suggest/completion/FuzzyOptions.java +++ b/src/main/java/org/codelibs/fesen/opensearch/search/suggest/completion/FuzzyOptions.java @@ -43,6 +43,7 @@ import org.codelibs.fesen.opensearch.core.xcontent.ToXContentFragment; import org.codelibs.fesen.opensearch.core.xcontent.XContentBuilder; import org.codelibs.fesen.opensearch.core.xcontent.XContentParser; +import org.codelibs.fesen.opensearch.index.query.RegexpQueryBuilder; import java.io.IOException; import java.util.Objects; @@ -107,6 +108,31 @@ private FuzzyOptions( this.maxDeterminizedStates = maxDeterminizedStates; } + /** + * Validates {@code max_determinized_states} against the shared ceiling and returns it. An + * unbounded value disables Lucene's determinize safeguard and lets a crafted expansion exhaust the + * heap before Lucene throws its own complexity exception. See CVE-2026-63136 and + * {@link RegexpQueryBuilder#MAX_DETERMINIZE_WORK_LIMIT}. + * + * @param maxDeterminizedStates the value to validate + * @return the validated value + */ + private static int validateMaxDeterminizedStates(int maxDeterminizedStates) { + if (maxDeterminizedStates < 0) { + throw new IllegalArgumentException("maxDeterminizedStates must not be negative"); + } + if (maxDeterminizedStates > RegexpQueryBuilder.MAX_DETERMINIZE_WORK_LIMIT) { + throw new IllegalArgumentException( + "maxDeterminizedStates cannot exceed [" + + RegexpQueryBuilder.MAX_DETERMINIZE_WORK_LIMIT + + "] but was [" + + maxDeterminizedStates + + "]" + ); + } + return maxDeterminizedStates; + } + @Override public void writeTo(StreamOutput out) throws IOException { out.writeBoolean(transpositions); @@ -253,10 +279,7 @@ public Builder setFuzzyPrefixLength(int fuzzyPrefixLength) { * @return this instance */ public Builder setMaxDeterminizedStates(int maxDeterminizedStates) { - if (maxDeterminizedStates < 0) { - throw new IllegalArgumentException("maxDeterminizedStates must not be negative"); - } - this.maxDeterminizedStates = maxDeterminizedStates; + this.maxDeterminizedStates = validateMaxDeterminizedStates(maxDeterminizedStates); return this; } diff --git a/src/main/java/org/codelibs/fesen/opensearch/search/suggest/completion/RegexOptions.java b/src/main/java/org/codelibs/fesen/opensearch/search/suggest/completion/RegexOptions.java index 78ca85ed..040c4d21 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/search/suggest/completion/RegexOptions.java +++ b/src/main/java/org/codelibs/fesen/opensearch/search/suggest/completion/RegexOptions.java @@ -44,6 +44,7 @@ import org.codelibs.fesen.opensearch.core.xcontent.XContentBuilder; import org.codelibs.fesen.opensearch.core.xcontent.XContentParser; import org.codelibs.fesen.opensearch.index.query.RegexpFlag; +import org.codelibs.fesen.opensearch.index.query.RegexpQueryBuilder; import java.io.IOException; @@ -92,6 +93,31 @@ private RegexOptions(int flagsValue, int maxDeterminizedStates) { this.maxDeterminizedStates = maxDeterminizedStates; } + /** + * Validates {@code max_determinized_states} against the shared ceiling and returns it. An + * unbounded value disables Lucene's determinize safeguard and lets a crafted regex exhaust the + * heap before Lucene throws its own complexity exception. See CVE-2026-63136 and + * {@link RegexpQueryBuilder#MAX_DETERMINIZE_WORK_LIMIT}. + * + * @param maxDeterminizedStates the value to validate + * @return the validated value + */ + private static int validateMaxDeterminizedStates(int maxDeterminizedStates) { + if (maxDeterminizedStates < 0) { + throw new IllegalArgumentException("maxDeterminizedStates must not be negative"); + } + if (maxDeterminizedStates > RegexpQueryBuilder.MAX_DETERMINIZE_WORK_LIMIT) { + throw new IllegalArgumentException( + "maxDeterminizedStates cannot exceed [" + + RegexpQueryBuilder.MAX_DETERMINIZE_WORK_LIMIT + + "] but was [" + + maxDeterminizedStates + + "]" + ); + } + return maxDeterminizedStates; + } + @Override public void writeTo(StreamOutput out) throws IOException { out.writeVInt(flagsValue); @@ -164,10 +190,7 @@ private Builder setFlagsValue(int flagsValue) { * @return this instance */ public Builder setMaxDeterminizedStates(int maxDeterminizedStates) { - if (maxDeterminizedStates < 0) { - throw new IllegalArgumentException("maxDeterminizedStates must not be negative"); - } - this.maxDeterminizedStates = maxDeterminizedStates; + this.maxDeterminizedStates = validateMaxDeterminizedStates(maxDeterminizedStates); return this; } diff --git a/src/main/java/org/codelibs/fesen/opensearch/threadpool/ThreadPool.java b/src/main/java/org/codelibs/fesen/opensearch/threadpool/ThreadPool.java index 114a6df5..c918e368 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/threadpool/ThreadPool.java +++ b/src/main/java/org/codelibs/fesen/opensearch/threadpool/ThreadPool.java @@ -71,7 +71,6 @@ import java.util.concurrent.RejectedExecutionHandler; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledThreadPoolExecutor; -import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; @@ -254,7 +253,11 @@ public enum ThreadPoolType { /** * The FORK_JOIN value. */ - FORK_JOIN("fork_join"); + FORK_JOIN("fork_join"), + /** + * The VIRTUAL value. + */ + VIRTUAL("virtual"); private final String type; @@ -757,7 +760,10 @@ static class ExecutorHolder { public final Info info; ExecutorHolder(ExecutorService executor, Info info) { - assert executor instanceof OpenSearchThreadPoolExecutor || executor == DIRECT_EXECUTOR || executor instanceof ForkJoinPool; + assert executor instanceof OpenSearchThreadPoolExecutor + || executor == DIRECT_EXECUTOR + || executor instanceof ForkJoinPool + || info.type == ThreadPoolType.VIRTUAL; this.executor = executor; this.info = info; } @@ -865,12 +871,15 @@ public Info(StreamInput in) throws IOException { public void writeTo(StreamOutput out) throws IOException { out.writeString(name); if (type == ThreadPoolType.RESIZABLE && out.getVersion().before(Version.V_3_0_0)) { - // Opensearch on older version doesn't know about "resizable" thread pool. Convert RESIZABLE to FIXED + // OpenSearch on older version doesn't know about "resizable" thread pool. Convert RESIZABLE to FIXED // to avoid serialization/de-serization issue between nodes with different OpenSearch version out.writeString(ThreadPoolType.FIXED.getType()); } else if (type == ThreadPoolType.FORK_JOIN && out.getVersion().before(Version.V_3_4_0)) { - // Opensearch on older version doesn't know about "fork_join" thread pool. Convert FORK_JOIN to FIXED + // OpenSearch on older version doesn't know about "fork_join" thread pool. Convert FORK_JOIN to FIXED out.writeString(ThreadPoolType.FIXED.getType()); + } else if (type == ThreadPoolType.VIRTUAL && out.getVersion().before(Version.V_3_8_0)) { + // VIRTUAL introduced in 3.8, convert to SCALING for bwc + out.writeString(ThreadPoolType.SCALING.getType()); } else { out.writeString(type.getType()); } @@ -947,6 +956,8 @@ public XContentBuilder toXContent(XContentBuilder builder, Params params) throws } } else if (type == ThreadPoolType.FORK_JOIN) { builder.field("parallelism", max); + } else if (type == ThreadPoolType.VIRTUAL) { + // an unbounded virtual thread-per-task pool has no size, keep alive, or queue to report } else { assert max != -1; builder.field("size", max); diff --git a/src/main/java/org/codelibs/fesen/opensearch/transport/client/ClusterAdminClient.java b/src/main/java/org/codelibs/fesen/opensearch/transport/client/ClusterAdminClient.java index 918c30a1..9f6b973e 100644 --- a/src/main/java/org/codelibs/fesen/opensearch/transport/client/ClusterAdminClient.java +++ b/src/main/java/org/codelibs/fesen/opensearch/transport/client/ClusterAdminClient.java @@ -50,6 +50,8 @@ import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksRequest; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksRequestBuilder; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.cancel.CancelTasksResponse; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskRequest; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskRequestBuilder; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskRequest; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskRequestBuilder; import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.get.GetTaskResponse; @@ -462,6 +464,38 @@ public interface ClusterAdminClient extends OpenSearchClient { */ GetTaskRequestBuilder prepareGetTask(TaskId taskId); + /** + * Delete a stored completed task result. + * + * @param request the request + * @return The result future + */ + ActionFuture deleteTask(DeleteTaskRequest request); + + /** + * Delete a stored completed task result. + * + * @param request the request + * @param listener A listener to be notified with the result + */ + void deleteTask(DeleteTaskRequest request, ActionListener listener); + + /** + * Delete a stored completed task result by id. + * + * @param taskId the id of the task whose stored result is deleted + * @return the request builder + */ + DeleteTaskRequestBuilder prepareDeleteTask(String taskId); + + /** + * Delete a stored completed task result by id. + * + * @param taskId the id of the task whose stored result is deleted + * @return the request builder + */ + DeleteTaskRequestBuilder prepareDeleteTask(TaskId taskId); + /** * Cancel tasks * diff --git a/src/test/java/org/codelibs/fesen/client/OpenSearch3ClientTest.java b/src/test/java/org/codelibs/fesen/client/OpenSearch3ClientTest.java index 3c6feda7..60aacd1a 100644 --- a/src/test/java/org/codelibs/fesen/client/OpenSearch3ClientTest.java +++ b/src/test/java/org/codelibs/fesen/client/OpenSearch3ClientTest.java @@ -18,6 +18,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; import static org.codelibs.fesen.opensearch.core.action.ActionListener.wrap; @@ -129,7 +130,7 @@ class OpenSearch3ClientTest { static final Logger logger = Logger.getLogger(OpenSearch3ClientTest.class.getName()); - static final String version = "3.8.0"; + static final String version = "3.9.0"; static final String imageTag = "public.ecr.aws/opensearchproject/opensearch:" + version; @@ -2780,6 +2781,42 @@ void test_get_task() throws Exception { } } + @Test + void test_delete_task() throws Exception { + // A task started with wait_for_completion=false leaves its result in the .tasks index, which + // is the only thing the delete task API removes. + final String srcIndex = "test_delete_task_src"; + client.admin().indices().prepareCreate(srcIndex).execute().actionGet(); + client.prepareIndex().setIndex(srcIndex).setId("1").setRefreshPolicy(RefreshPolicy.IMMEDIATE) + .setSource("{\"text\":\"test\"}", XContentType.JSON).execute().actionGet(); + + final String url = "http://" + server.getHost() + ":" + server.getFirstMappedPort(); + final String started; + try (CurlResponse response = Curl.post(url + "/_reindex?wait_for_completion=false").header("Content-Type", "application/json") + .body("{\"source\":{\"index\":\"" + srcIndex + "\"},\"dest\":{\"index\":\"test_delete_task_dest\"}}").execute()) { + started = response.getContentAsString(); + } + final java.util.regex.Matcher matcher = java.util.regex.Pattern.compile("\"task\"\\s*:\\s*\"([^\"]+)\"").matcher(started); + assertTrue(matcher.find(), started); + final TaskId taskId = new TaskId(matcher.group(1)); + + // Wait for the task to finish so that its result is stored. + boolean completed = false; + for (int i = 0; i < 50 && !completed; i++) { + completed = + client.admin().cluster().prepareGetTask(taskId).execute().actionGet().getTask().getResponseAsMap().isEmpty() == false; + if (!completed) { + Thread.sleep(200L); + } + } + assertTrue(completed); + + assertTrue(client.admin().cluster().prepareDeleteTask(taskId.toString()).execute().actionGet().isAcknowledged()); + + // The stored result is gone. + assertThrows(OpenSearchException.class, () -> client.admin().cluster().prepareGetTask(taskId).execute().actionGet()); + } + @Test void test_upgrade_status() throws Exception { final String index = "test_upgrade_status"; diff --git a/src/test/java/org/codelibs/fesen/client/action/ActionTestUtils.java b/src/test/java/org/codelibs/fesen/client/action/ActionTestUtils.java index 0a8a3cf5..ce8207ef 100644 --- a/src/test/java/org/codelibs/fesen/client/action/ActionTestUtils.java +++ b/src/test/java/org/codelibs/fesen/client/action/ActionTestUtils.java @@ -175,6 +175,23 @@ static String url(final CurlRequest request) { } } + /** + * Returns the HTTP method of the given curl request, read reflectively from the + * protected {@code method} field. + * + * @param request the curl request to inspect + * @return the HTTP method name + */ + static String method(final CurlRequest request) { + try { + final Field field = CurlRequest.class.getDeclaredField("method"); + field.setAccessible(true); + return String.valueOf(field.get(request)); + } catch (final ReflectiveOperationException e) { + throw new RuntimeException(e); + } + } + private static String decode(final String value) { return URLDecoder.decode(value, StandardCharsets.UTF_8); } diff --git a/src/test/java/org/codelibs/fesen/client/action/HttpDeleteTaskActionTest.java b/src/test/java/org/codelibs/fesen/client/action/HttpDeleteTaskActionTest.java new file mode 100644 index 00000000..558a1883 --- /dev/null +++ b/src/test/java/org/codelibs/fesen/client/action/HttpDeleteTaskActionTest.java @@ -0,0 +1,69 @@ +/* + * Copyright 2012-2025 CodeLibs Project and the Others. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, + * either express or implied. See the License for the specific language + * governing permissions and limitations under the License. + */ +package org.codelibs.fesen.client.action; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; + +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskAction; +import org.codelibs.fesen.opensearch.action.admin.cluster.node.tasks.delete.DeleteTaskRequest; +import org.codelibs.fesen.opensearch.action.support.clustermanager.AcknowledgedResponse; +import org.codelibs.fesen.opensearch.common.xcontent.json.JsonXContent; +import org.codelibs.fesen.opensearch.core.tasks.TaskId; +import org.codelibs.fesen.opensearch.core.xcontent.DeprecationHandler; +import org.codelibs.fesen.opensearch.core.xcontent.NamedXContentRegistry; +import org.codelibs.fesen.opensearch.core.xcontent.XContentParser; +import org.junit.jupiter.api.Test; + +class HttpDeleteTaskActionTest { + + private final HttpDeleteTaskAction clientAction = new HttpDeleteTaskAction(ActionTestUtils.testClient(), DeleteTaskAction.INSTANCE); + + @Test + void test_getCurlRequest_deletesTaskById() { + final DeleteTaskRequest request = new DeleteTaskRequest().setTaskId(new TaskId("node1:42")); + final org.codelibs.curl.CurlRequest curlRequest = clientAction.getCurlRequest(request); + assertEquals("DELETE", ActionTestUtils.method(curlRequest)); + assertTrue(ActionTestUtils.url(curlRequest).endsWith("/_tasks/node1:42"), ActionTestUtils.url(curlRequest)); + } + + @Test + void test_request_requiresTaskId() { + assertNotNull(new DeleteTaskRequest().validate()); + assertEquals(null, new DeleteTaskRequest().setTaskId(new TaskId("node1:42")).validate()); + } + + @Test + void test_fromXContent_acknowledged() throws IOException { + try (final XContentParser parser = JsonXContent.jsonXContent.createParser(NamedXContentRegistry.EMPTY, + DeprecationHandler.THROW_UNSUPPORTED_OPERATION, "{\"acknowledged\":true}")) { + assertTrue(AcknowledgedResponse.fromXContent(parser).isAcknowledged()); + } + } + + @Test + void test_fromXContent_errorBody_throws() throws IOException { + // A non-2xx error body is routed through the success path; the parser must reject it. + try (final XContentParser parser = JsonXContent.jsonXContent.createParser(NamedXContentRegistry.EMPTY, + DeprecationHandler.THROW_UNSUPPORTED_OPERATION, "{\"error\":{\"type\":\"resource_not_found_exception\"},\"status\":404}")) { + assertThrows(Exception.class, () -> AcknowledgedResponse.fromXContent(parser)); + } + } +} diff --git a/src/test/java/org/codelibs/fesen/client/action/HttpNodesStatsActionTest.java b/src/test/java/org/codelibs/fesen/client/action/HttpNodesStatsActionTest.java index f23299b8..0f5fa5c0 100644 --- a/src/test/java/org/codelibs/fesen/client/action/HttpNodesStatsActionTest.java +++ b/src/test/java/org/codelibs/fesen/client/action/HttpNodesStatsActionTest.java @@ -37,6 +37,7 @@ import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; +import org.codelibs.fesen.opensearch.action.ActionConcurrencyLimiterStats; import org.codelibs.fesen.opensearch.action.admin.cluster.node.stats.NodeStats; import org.codelibs.fesen.opensearch.action.admin.cluster.node.stats.NodesStatsAction; import org.codelibs.fesen.opensearch.action.admin.cluster.node.stats.NodesStatsRequest; @@ -45,6 +46,8 @@ import org.codelibs.fesen.opensearch.common.xcontent.json.JsonXContent; import org.codelibs.fesen.opensearch.core.xcontent.DeprecationHandler; import org.codelibs.fesen.opensearch.core.xcontent.NamedXContentRegistry; +import org.codelibs.fesen.opensearch.core.xcontent.ToXContent; +import org.codelibs.fesen.opensearch.core.xcontent.XContentBuilder; import org.codelibs.fesen.opensearch.core.xcontent.XContentParser; import org.codelibs.fesen.opensearch.plugin.stats.NativeAllocatorPoolStats; import org.codelibs.fesen.opensearch.plugins.BlockCacheStats; @@ -293,6 +296,88 @@ void test_parseNodeStats_withRemoteStore() throws Exception { } } + @Test + @Timeout(value = 5, unit = TimeUnit.SECONDS) + void test_parseNodeStats_withConcurrencyLimiters() throws Exception { + final String json = "{" + "\"name\":\"test-node\"," + "\"timestamp\":1234567890," + "\"concurrency_limiters\":{" + + "\"search\":{\"action_name\":\"indices:data/read/search\",\"mode\":\"enforced\",\"algorithm\":\"gradient\"," + + "\"current_limit\":20,\"in_flight\":3,\"total_rejected\":7,\"last_rtt_millis\":15,\"rtt_no_load_millis\":9}," + + "\"bulk\":{\"action_name\":\"indices:data/write/bulk\",\"mode\":\"monitor_only\",\"algorithm\":\"vegas\"," + + "\"current_limit\":50,\"in_flight\":0,\"total_rejected\":0,\"unknown\":{\"nested\":[1,2]}}" + "}," + + "\"transport_address\":\"127.0.0.1:9300\"" + "}"; + try (final XContentParser parser = createParser(json)) { + parser.nextToken(); + final NodeStats nodeStats = callParseNodeStats(parser, "node1"); + assertEquals("test-node", nodeStats.getNode().getName()); + final ActionConcurrencyLimiterStats stats = nodeStats.getConcurrencyLimiterStats(); + assertNotNull(stats); + assertEquals(2, stats.getSnapshots().size()); + final ActionConcurrencyLimiterStats.ActionLimiterSnapshot search = stats.getSnapshots().get(0); + assertEquals("search", search.getAlias()); + assertEquals("indices:data/read/search", search.getActionName()); + assertEquals("enforced", search.getMode()); + assertEquals("gradient", search.getAlgorithm()); + assertEquals(20, search.getCurrentLimit()); + assertEquals(3, search.getInFlight()); + assertEquals(7L, search.getTotalRejected()); + assertEquals(15L, search.getLastRttMillis()); + assertEquals(9L, search.getRttNoLoadMillis()); + final ActionConcurrencyLimiterStats.ActionLimiterSnapshot bulk = stats.getSnapshots().get(1); + assertEquals("bulk", bulk.getAlias()); + assertEquals(50, bulk.getCurrentLimit()); + // The server omits both round-trip times while they are unavailable; the sentinel is -1. + assertEquals(-1L, bulk.getLastRttMillis()); + assertEquals(-1L, bulk.getRttNoLoadMillis()); + } + } + + @Test + @Timeout(value = 5, unit = TimeUnit.SECONDS) + void test_parseNodeStats_withoutConcurrencyLimiters() throws Exception { + final String json = + "{" + "\"name\":\"test-node\"," + "\"timestamp\":1234567890," + "\"transport_address\":\"127.0.0.1:9300\"" + "}"; + try (final XContentParser parser = createParser(json)) { + parser.nextToken(); + assertNull(callParseNodeStats(parser, "node1").getConcurrencyLimiterStats()); + } + } + + @Test + @Timeout(value = 5, unit = TimeUnit.SECONDS) + void test_parseConcurrencyLimiters_roundTripsServerRendering() throws Exception { + final ActionConcurrencyLimiterStats expected = new ActionConcurrencyLimiterStats(List.of( + new ActionConcurrencyLimiterStats.ActionLimiterSnapshot("search", "indices:data/read/search", "enforced", "gradient", 20, 3, + 7L, 15L, 9L), + new ActionConcurrencyLimiterStats.ActionLimiterSnapshot("bulk", "indices:data/write/bulk", "monitor_only", "vegas", 50, 0, + 0L, -1L, -1L))); + final XContentBuilder builder = JsonXContent.contentBuilder().startObject(); + expected.toXContent(builder, ToXContent.EMPTY_PARAMS); + builder.endObject(); + try (final XContentParser parser = createParser(builder.toString())) { + // createParser has already advanced to the outer START_OBJECT + parser.nextToken(); // FIELD_NAME concurrency_limiters + parser.nextToken(); // START_OBJECT + parser.nextToken(); // first alias + final Method m = HttpNodesStatsAction.class.getDeclaredMethod("parseConcurrencyLimiterStats", XContentParser.class); + m.setAccessible(true); + final ActionConcurrencyLimiterStats actual = (ActionConcurrencyLimiterStats) m.invoke(action, parser); + assertEquals(2, actual.getSnapshots().size()); + for (int i = 0; i < 2; i++) { + final ActionConcurrencyLimiterStats.ActionLimiterSnapshot e = expected.getSnapshots().get(i); + final ActionConcurrencyLimiterStats.ActionLimiterSnapshot a = actual.getSnapshots().get(i); + assertEquals(e.getAlias(), a.getAlias()); + assertEquals(e.getActionName(), a.getActionName()); + assertEquals(e.getMode(), a.getMode()); + assertEquals(e.getAlgorithm(), a.getAlgorithm()); + assertEquals(e.getCurrentLimit(), a.getCurrentLimit()); + assertEquals(e.getInFlight(), a.getInFlight()); + assertEquals(e.getTotalRejected(), a.getTotalRejected()); + assertEquals(e.getLastRttMillis(), a.getLastRttMillis()); + assertEquals(e.getRttNoLoadMillis(), a.getRttNoLoadMillis()); + } + } + } + @Test @Timeout(value = 5, unit = TimeUnit.SECONDS) void test_parseNodeStats_withRoles() throws Exception {