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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
8 changes: 4 additions & 4 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
<artifactId>fesen-httpclient</artifactId>
<packaging>jar</packaging>
<name>fesen-httpclient</name>
<version>3.8.1-SNAPSHOT</version>
<version>3.9.0-SNAPSHOT</version>
<description>HTTP client for OpenSearch</description>
<url>https://github.com/codelibs/fesen-httpclient</url>
<inceptionYear>2012</inceptionYear>
Expand Down Expand Up @@ -38,9 +38,9 @@
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<maven.compiler.release>21</maven.compiler.release>
<lucene.version>10.5.0</lucene.version>
<jackson.version>3.2.1</jackson.version>
<log4j.version>2.25.4</log4j.version>
<lucene.version>10.5.1</lucene.version>
<jackson.version>3.2.2</jackson.version>
<log4j.version>2.25.5</log4j.version>
<junit.jupiter.version>5.12.2</junit.jupiter.version>
<testcontainers.version>1.21.4</testcontainers.version>
</properties>
Expand Down
23 changes: 23 additions & 0 deletions src/main/java/org/codelibs/fesen/client/HttpAbstractClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -936,6 +939,26 @@ public GetTaskRequestBuilder prepareGetTask(final TaskId taskId) {
return new GetTaskRequestBuilder(this, GetTaskAction.INSTANCE).setTaskId(taskId);
}

@Override
public ActionFuture<AcknowledgedResponse> deleteTask(final DeleteTaskRequest request) {
return execute(DeleteTaskAction.INSTANCE, request);
}

@Override
public void deleteTask(final DeleteTaskRequest request, final ActionListener<AcknowledgedResponse> 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<CancelTasksResponse> cancelTasks(final CancelTasksRequest request) {
return execute(CancelTasksAction.INSTANCE, request);
Expand Down
8 changes: 8 additions & 0 deletions src/main/java/org/codelibs/fesen/client/HttpClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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<AcknowledgedResponse> actionListener = (ActionListener<AcknowledgedResponse>) listener;
new HttpDeleteTaskAction(this, DeleteTaskAction.INSTANCE).execute((DeleteTaskRequest) request, actionListener);
});
actions.put(GetTaskAction.INSTANCE, (request, listener) -> {
@SuppressWarnings("unchecked")
final ActionListener<GetTaskResponse> actionListener = (ActionListener<GetTaskResponse>) listener;
Expand Down
Original file line number Diff line number Diff line change
@@ -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<AcknowledgedResponse> 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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<String, String> attributes = new HashMap<>();
XContentParser.Token token;
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -438,6 +442,7 @@ protected NodeStats parseNodeStats(final XContentParser parser, final String nod
nodeCacheStats, //
remoteStoreNodeStats, //
nativeAllocatorStats, //
concurrencyLimiterStats, //
totalEstimatedNativeBytes);
}

Expand Down Expand Up @@ -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<ActionConcurrencyLimiterStats.ActionLimiterSnapshot> 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.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
};
Expand Down
10 changes: 9 additions & 1 deletion src/main/java/org/codelibs/fesen/opensearch/Version.java
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading