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);