diff --git a/pom.xml b/pom.xml
index d364d2ca..46784426 100644
--- a/pom.xml
+++ b/pom.xml
@@ -79,6 +79,7 @@
${project.basedir}/src/test/java/
UTF-8
UTF-8
+ 12.1.10
2.0.2
3.3.2
@@ -99,7 +100,17 @@
org.eclipse.jetty
jetty-http
- 9.4.52.v20230823
+ ${jetty.version}
+
+
+ org.eclipse.jetty
+ jetty-client
+ ${jetty.version}
+
+
+ org.eclipse.jetty
+ jetty-util
+ ${jetty.version}
io.cdap.cdap
@@ -270,6 +281,22 @@
org.eclipse.jetty
jetty-http
+
+ org.eclipse.jetty
+ jetty-client
+
+
+ org.eclipse.jetty
+ jetty-util
+
+
+ org.eclipse.jetty
+ jetty-util-ajax
+
+
+ org.eclipse.jetty
+ jetty-io
+
@@ -377,6 +404,14 @@
log4j
log4j
+
+ org.eclipse.jetty
+ *
+
+
+ org.eclipse.jetty.orbit
+ *
+
test
diff --git a/src/main/java/io/cdap/plugin/salesforce/SalesforceQueryUtil.java b/src/main/java/io/cdap/plugin/salesforce/SalesforceQueryUtil.java
index cfbb0cfd..d3847ff4 100644
--- a/src/main/java/io/cdap/plugin/salesforce/SalesforceQueryUtil.java
+++ b/src/main/java/io/cdap/plugin/salesforce/SalesforceQueryUtil.java
@@ -21,12 +21,11 @@
import io.cdap.plugin.salesforce.authenticator.AuthenticatorCredentials;
import io.cdap.plugin.salesforce.parser.SalesforceQueryParser;
import io.cdap.plugin.salesforce.plugin.OAuthInfo;
+import org.eclipse.jetty.client.ContentResponse;
import org.eclipse.jetty.client.HttpClient;
-import org.eclipse.jetty.client.api.ContentResponse;
-import org.eclipse.jetty.client.api.Request;
+import org.eclipse.jetty.client.Request;
import org.eclipse.jetty.http.HttpHeader;
import org.eclipse.jetty.http.HttpMethod;
-import org.eclipse.jetty.util.ssl.SslContextFactory;
import java.io.IOException;
import java.net.HttpURLConnection;
@@ -152,8 +151,7 @@ public static QueryPlanResponse getQueryPlan(
SalesforceConstants.API_VERSION,
URLEncoder.encode(query, "UTF-8"));
- SslContextFactory sslContextFactory = new SslContextFactory();
- HttpClient httpClient = new HttpClient(sslContextFactory);
+ HttpClient httpClient = new HttpClient();
httpClient.setConnectTimeout(credentials.getConnectTimeout());
if (!Strings.isNullOrEmpty(credentials.getProxyUrl())) {
Authenticator.setProxy(credentials, httpClient);
@@ -163,12 +161,9 @@ public static QueryPlanResponse getQueryPlan(
httpClient.start();
Request request = httpClient.newRequest(explainUrl)
.method(HttpMethod.GET)
- .header(
- HttpHeader.AUTHORIZATION,
- "Bearer " + oAuthInfo.getAccessToken())
- .header(
- HttpHeader.CONTENT_TYPE,
- "application/json");
+ .headers(headers -> headers
+ .put(HttpHeader.AUTHORIZATION, "Bearer " + oAuthInfo.getAccessToken())
+ .put(HttpHeader.CONTENT_TYPE, "application/json"));
ContentResponse response = request.send();
String responseContent = response.getContentAsString();
if (response.getStatus() != HttpURLConnection.HTTP_OK) {
diff --git a/src/main/java/io/cdap/plugin/salesforce/authenticator/Authenticator.java b/src/main/java/io/cdap/plugin/salesforce/authenticator/Authenticator.java
index 8bbf280c..c4c87c2d 100644
--- a/src/main/java/io/cdap/plugin/salesforce/authenticator/Authenticator.java
+++ b/src/main/java/io/cdap/plugin/salesforce/authenticator/Authenticator.java
@@ -25,8 +25,7 @@
import org.eclipse.jetty.client.HttpClient;
import org.eclipse.jetty.client.HttpProxy;
import org.eclipse.jetty.client.ProxyConfiguration;
-import org.eclipse.jetty.client.api.Request;
-import org.eclipse.jetty.util.ssl.SslContextFactory;
+import org.eclipse.jetty.client.Request;
import java.net.URI;
import java.net.URISyntaxException;
@@ -89,8 +88,7 @@ public static OAuthInfo getOAuthInfo(AuthenticatorCredentials credentials) throw
throw new IllegalArgumentException("Grant type cannot be null for OAuth flow to fetch access token.");
}
- SslContextFactory sslContextFactory = new SslContextFactory();
- HttpClient httpClient = new HttpClient(sslContextFactory);
+ HttpClient httpClient = new HttpClient();
httpClient.setConnectTimeout(credentials.getConnectTimeout());
if (!Strings.isNullOrEmpty(credentials.getProxyUrl())) {
setProxy(credentials, httpClient);
diff --git a/src/main/java/io/cdap/plugin/salesforce/plugin/source/streaming/SalesforcePushTopicListener.java b/src/main/java/io/cdap/plugin/salesforce/plugin/source/streaming/SalesforcePushTopicListener.java
index ef7f34da..2bfa64cb 100644
--- a/src/main/java/io/cdap/plugin/salesforce/plugin/source/streaming/SalesforcePushTopicListener.java
+++ b/src/main/java/io/cdap/plugin/salesforce/plugin/source/streaming/SalesforcePushTopicListener.java
@@ -32,8 +32,7 @@
import org.cometd.common.JSONContext;
import org.cometd.common.JacksonJSONContextClient;
import org.eclipse.jetty.client.HttpClient;
-import org.eclipse.jetty.client.api.Request;
-import org.eclipse.jetty.util.ssl.SslContextFactory;
+import org.eclipse.jetty.client.Request;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -129,10 +128,8 @@ public String getMessage(long timeout, TimeUnit unit) throws InterruptedExceptio
private BayeuxClient getClient(AuthenticatorCredentials credentials) throws Exception {
OAuthInfo oAuthInfo = Authenticator.getOAuthInfo(credentials);
- SslContextFactory sslContextFactory = new SslContextFactory();
-
// Set up a Jetty HTTP client to use with CometD
- HttpClient httpClient = new HttpClient(sslContextFactory);
+ HttpClient httpClient = new HttpClient();
httpClient.setConnectTimeout(CONNECTION_TIMEOUT_MS);
if (!Strings.isNullOrEmpty(credentials.getProxyUrl())) {
Authenticator.setProxy(credentials, httpClient);
@@ -146,15 +143,15 @@ private BayeuxClient getClient(AuthenticatorCredentials credentials) throws Exce
Map transportOptions = new HashMap<>();
transportOptions.put(ClientTransport.JSON_CONTEXT_OPTION, jsonContext);
- // Adds the OAuth header in LongPollingTransport
- LongPollingTransport transport = new LongPollingTransport(
- transportOptions, httpClient) {
+ // Adds the OAuth header for all CometD requests
+ httpClient.getRequestListeners().addListener(new Request.Listener() {
@Override
- protected void customize(Request exchange) {
- super.customize(exchange);
- exchange.header("Authorization", "OAuth " + oAuthInfo.getAccessToken());
+ public void onBegin(Request request) {
+ request.headers(headers -> headers.put("Authorization", "OAuth " + oAuthInfo.getAccessToken()));
}
- };
+ });
+
+ LongPollingTransport transport = new LongPollingTransport(transportOptions, httpClient);
// Now set up the Bayeux client itself
return new BayeuxClient(oAuthInfo.getInstanceURL() + DEFAULT_PUSH_ENDPOINT, transport);
diff --git a/src/test/java/io/cdap/plugin/salesforce/SalesforceQueryUtilTest.java b/src/test/java/io/cdap/plugin/salesforce/SalesforceQueryUtilTest.java
index f8aec01a..cb53de56 100644
--- a/src/test/java/io/cdap/plugin/salesforce/SalesforceQueryUtilTest.java
+++ b/src/test/java/io/cdap/plugin/salesforce/SalesforceQueryUtilTest.java
@@ -20,12 +20,11 @@
import io.cdap.plugin.salesforce.authenticator.AuthenticatorCredentials;
import io.cdap.plugin.salesforce.plugin.OAuthInfo;
import io.cdap.plugin.salesforce.plugin.source.batch.util.SalesforceSplitUtil;
+import org.eclipse.jetty.client.ContentResponse;
import org.eclipse.jetty.client.HttpClient;
-import org.eclipse.jetty.client.api.ContentResponse;
-import org.eclipse.jetty.client.api.Request;
+import org.eclipse.jetty.client.Request;
import org.eclipse.jetty.http.HttpHeader;
import org.eclipse.jetty.http.HttpMethod;
-import org.eclipse.jetty.util.ssl.SslContextFactory;
import org.junit.Assert;
import org.junit.Test;
import org.junit.runner.RunWith;
@@ -296,13 +295,12 @@ public void getQueryPlan_success_returnsQueryPlanResponse() throws Exception {
PowerMockito.mockStatic(Authenticator.class);
PowerMockito.when(Authenticator.getOAuthInfo(credentials)).thenReturn(oAuthInfo);
HttpClient httpClient = PowerMockito.mock(HttpClient.class);
- PowerMockito.whenNew(HttpClient.class).withArguments(Mockito.any(SslContextFactory.class))
+ PowerMockito.whenNew(HttpClient.class).withNoArguments()
.thenReturn(httpClient);
Request request = Mockito.mock(Request.class);
Mockito.when(httpClient.newRequest(Mockito.anyString())).thenReturn(request);
Mockito.when(request.method(Mockito.any(HttpMethod.class))).thenReturn(request);
- Mockito.when(request.header(Mockito.any(HttpHeader.class), Mockito.anyString()))
- .thenReturn(request);
+ Mockito.when(request.headers(Mockito.any())).thenReturn(request);
ContentResponse response = Mockito.mock(ContentResponse.class);
Mockito.when(request.send()).thenReturn(response);
Mockito.when(response.getStatus()).thenReturn(HttpURLConnection.HTTP_OK);