Skip to content

Commit aeb6b11

Browse files
authored
Merge pull request #125 from forcedotcom/develop
@W-12387797: 1.17.0 release for tableau user to use hyper engine as default option
2 parents 5fc9c93 + 2920b65 commit aeb6b11

13 files changed

Lines changed: 154 additions & 164 deletions

File tree

README.md

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,9 @@ Class.forName("com.salesforce.cdp.queryservice.QueryServiceDriver");
119119
Properties properties = new Properties();
120120
properties.put("user", <UserName>);
121121
properties.put("password", <Password>);
122+
properties.put("clientId", <Client Id of the connected App>);
123+
properties.put("clientSecret", <Client Secret of the connected App>);
124+
122125
123126
Connection connection = DriverManager.getConnection("jdbc:queryService-jdbc:https://login.salesforce.com", properties);
124127
```
@@ -164,7 +167,9 @@ import jaydebeapi
164167
// Sample properties with username and password flow.
165168
properties = {
166169
'user': "<UserName>",
167-
'password': "<Password>"
170+
'password': "<Password>",
171+
'clientId': "<Client Id of the connected App>",
172+
'clientSecret': "<Client Secret of the connected App>"
168173
}
169174
170175
// Sample properties with key-pair authentication flow.

pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66

77
<groupId>com.queryService</groupId>
88
<artifactId>Salesforce-CDP-jdbc</artifactId>
9-
<version>1.16.0</version>
9+
<version>1.17.0</version>
1010

1111
<properties>
1212
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>

src/main/java/com/salesforce/cdp/queryservice/core/QueryServiceAbstractStatement.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
import com.google.protobuf.Struct;
2020
import com.google.protobuf.Value;
2121
import com.salesforce.a360.queryservice.grpc.v1.AnsiSqlQueryStreamResponse;
22+
import com.salesforce.cdp.queryservice.enums.QueryEngineEnum;
2223
import com.salesforce.cdp.queryservice.model.QueryServiceResponse;
2324
import com.salesforce.cdp.queryservice.model.Type;
2425
import com.salesforce.cdp.queryservice.util.ArrowUtil;
@@ -80,15 +81,15 @@ public QueryServiceAbstractStatement(QueryServiceConnection queryServiceConnecti
8081
public ResultSet executeQuery(String sql) throws SQLException {
8182
try {
8283
this.sql = sql;
83-
boolean isEnableStreamFlow = this.connection.isEnableStreamFlow();
84+
QueryEngineEnum engineEnum = this.connection.getQueryEngineEnum();
8485

8586
boolean isCursorBasedPaginationReq = this.connection.isCursorBasedPaginationReq();
8687

8788
boolean requireManagedPagination = isTableauQuery() && !isCursorBasedPaginationReq;
8889
Optional<Integer> limit = requireManagedPagination ? Optional.of(Constants.MAX_LIMIT) : Optional.empty();
8990
Optional<String> orderby = requireManagedPagination ? Optional.of("1 ASC") : Optional.empty();
9091

91-
if(isEnableStreamFlow) {
92+
if (QueryEngineEnum.HYPER == engineEnum) {
9293
Iterator<AnsiSqlQueryStreamResponse> response = queryGrpcExecutor.executeQueryWithRetry(sql);
9394
return createResultSetFromResponse(response);
9495
} else {

src/main/java/com/salesforce/cdp/queryservice/core/QueryServiceConnection.java

Lines changed: 48 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -17,22 +17,29 @@
1717
package com.salesforce.cdp.queryservice.core;
1818

1919
import com.google.common.annotations.VisibleForTesting;
20+
import com.salesforce.cdp.queryservice.enums.QueryEngineEnum;
21+
import com.salesforce.cdp.queryservice.model.QueryConfigResponse;
2022
import com.salesforce.cdp.queryservice.model.Token;
2123
import com.salesforce.cdp.queryservice.util.Constants;
24+
import com.salesforce.cdp.queryservice.util.HttpHelper;
25+
import com.salesforce.cdp.queryservice.util.QueryExecutor;
26+
import com.salesforce.cdp.queryservice.util.TokenHelper;
2227
import lombok.extern.slf4j.Slf4j;
28+
import okhttp3.Response;
2329
import org.apache.commons.lang3.StringUtils;
2430

31+
import java.io.IOException;
2532
import java.sql.*;
2633
import java.util.Map;
2734
import java.util.Properties;
2835
import java.util.concurrent.Executor;
2936
import java.util.concurrent.atomic.AtomicBoolean;
3037

38+
import static com.salesforce.cdp.queryservice.util.Messages.QUERY_CONFIG_ERROR;
39+
3140
@Slf4j
3241
public class QueryServiceConnection implements Connection {
3342

34-
private static final String TEST_CONNECT_QUERY = "select 1";
35-
3643
private AtomicBoolean closed = new AtomicBoolean(false);
3744
private Properties properties;
3845
private final String serviceRootUrl;
@@ -42,12 +49,15 @@ public class QueryServiceConnection implements Connection {
4249
private final boolean isSocksProxyDisabled;
4350
private boolean enableStreamFlow = false;
4451
private String tenantUrl;
52+
private QueryEngineEnum queryEngineEnum;
53+
54+
private boolean isValid = false;
4555

4656
public QueryServiceConnection(String url, Properties properties) throws SQLException {
4757
this.properties = properties; // fixme: do deepCopy and modify the props
4858
this.serviceRootUrl = getServiceRootUrl(url);
4959
this.properties.put(Constants.LOGIN_URL, serviceRootUrl);
50-
addClientSecretsIfRequired(serviceRootUrl, this.properties);
60+
addClientUsernameIfRequired(this.properties);
5161

5262
// default `enableArrowStream` is false
5363
enableArrowStream = Boolean.parseBoolean(this.properties.getProperty(Constants.ENABLE_ARROW_STREAM));
@@ -57,8 +67,12 @@ public QueryServiceConnection(String url, Properties properties) throws SQLExcep
5767

5868
this.isSocksProxyDisabled = Boolean.parseBoolean(this.properties.getProperty(Constants.DISABLE_SOCKS_PROXY));
5969

70+
boolean isTableauConnection = Constants.TABLEAU_USER_AGENT_VALUE.equals(properties.getProperty(Constants.USER_AGENT));
71+
6072
// default `enableStreamFlow` is false
61-
enableStreamFlow = Boolean.parseBoolean(this.properties.getProperty(Constants.ENABLE_STREAM_FLOW, Constants.FALSE_STR));
73+
enableStreamFlow = isTableauConnection || Boolean.parseBoolean(this.properties.getProperty(Constants.ENABLE_STREAM_FLOW, Constants.FALSE_STR));
74+
75+
log.info("isTableauConnection {}, enableStreamFlow {}", isTableauConnection, enableStreamFlow);
6276

6377
// use isValid to test connection
6478
this.isValid(20);
@@ -87,37 +101,14 @@ static String getServiceRootUrl(String url) throws SQLException {
87101
/**
88102
* Adds client secrets to properties if not present and service url matches one of the existing envs.
89103
*
90-
* @param serviceRootUrl service url which is used to infer the environment
91104
* @param properties Properties containing the config
92105
* @throws SQLException when given service url doesn't match any envs and config doesn't have exists secrets
93106
*/
94107
@VisibleForTesting
95-
static void addClientSecretsIfRequired(String serviceRootUrl, Properties properties) throws SQLException {
108+
static void addClientUsernameIfRequired(Properties properties) throws SQLException {
96109
if (properties.containsKey(Constants.USER) && !properties.containsKey(Constants.USER_NAME)) {
97110
properties.put(Constants.USER_NAME, properties.get(Constants.USER));
98111
}
99-
100-
if (properties.containsKey(Constants.USER_NAME)
101-
&& !properties.containsKey(Constants.CLIENT_ID)
102-
&& !properties.containsKey(Constants.CLIENT_SECRET)
103-
&& !properties.containsKey(Constants.PRIVATE_KEY)) {
104-
log.debug("adding client secrets for server {}", serviceRootUrl);
105-
String serverUrl = serviceRootUrl.toLowerCase();
106-
if (serverUrl.endsWith(Constants.NA45_SERVER_URL)) {
107-
properties.put(Constants.CLIENT_ID, Constants.NA45_DEFAULT_CLIENT_ID);
108-
properties.put(Constants.CLIENT_SECRET, Constants.NA45_DEFAULT_CLIENT_SECRET);
109-
} else if (serverUrl.endsWith(Constants.NA46_SERVER_URL)) {
110-
properties.put(Constants.CLIENT_ID, Constants.NA46_DEFAULT_CLIENT_ID);
111-
properties.put(Constants.CLIENT_SECRET, Constants.NA46_DEFAULT_CLIENT_SECRET);
112-
} else if (serverUrl.endsWith(Constants.PROD_SERVER_URL)) {
113-
properties.put(Constants.CLIENT_ID, Constants.PROD_DEFAULT_CLIENT_ID);
114-
properties.put(Constants.CLIENT_SECRET, Constants.PROD_DEFAULT_CLIENT_SECRET);
115-
} else {
116-
throw new SQLException("specified url didn't match any existing envs");
117-
}
118-
} else {
119-
log.debug("No client secrets added for server {}", serviceRootUrl);
120-
}
121112
}
122113

123114
public boolean getEnableArrowStream() {
@@ -141,6 +132,10 @@ public boolean updateStreamFlow(boolean flag) {
141132
return enableStreamFlow;
142133
}
143134

135+
public QueryEngineEnum getQueryEngineEnum() {
136+
return queryEngineEnum;
137+
}
138+
144139
@Override
145140
public Statement createStatement() throws SQLException {
146141
return createStatement(ResultSet.TYPE_FORWARD_ONLY,
@@ -351,17 +346,17 @@ public boolean isValid(int timeout) throws SQLException {
351346
}
352347

353348
try {
354-
PreparedStatement statement = this.prepareStatement(TEST_CONNECT_QUERY);
355-
return statement.execute();
349+
if(this.isValid) {
350+
log.info("Reusing connection");
351+
return true;
352+
}
353+
354+
QueryConfigResponse configResponse = getQueryConfigResponse();
355+
this.queryEngineEnum = QueryEngineEnum.fromValue(configResponse.getQueryengine());
356+
this.isValid = true;
357+
return true;
356358
} catch (Exception e) {
357359
log.error("Exception while connecting to server", e);
358-
if(isEnableStreamFlow()) {
359-
// use http v2 api if hyper gRPC call is failing
360-
updateStreamFlow(false);
361-
try(PreparedStatement statement = this.prepareStatement(TEST_CONNECT_QUERY)) {
362-
return statement.execute();
363-
}
364-
}
365360
throw e;
366361
}
367362
}
@@ -457,4 +452,20 @@ public String getTenantUrl() {
457452
public void setTenantUrl(String tenantUrl) {
458453
this.tenantUrl = tenantUrl;
459454
}
455+
456+
private QueryExecutor createQueryExecutor() {
457+
return new QueryExecutor(this);
458+
}
459+
460+
QueryConfigResponse getQueryConfigResponse() throws SQLException {
461+
try {
462+
QueryExecutor executor = createQueryExecutor();
463+
Response response = executor.getQueryConfig();
464+
465+
return HttpHelper.handleSuccessResponse(response, QueryConfigResponse.class, false);
466+
} catch (IOException e) {
467+
log.error("Exception while getting config from query service", e);
468+
throw new SQLException(QUERY_CONFIG_ERROR, e);
469+
}
470+
}
460471
}
Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
package com.salesforce.cdp.queryservice.enums;
2+
3+
import com.fasterxml.jackson.annotation.JsonCreator;
4+
import com.fasterxml.jackson.annotation.JsonValue;
5+
6+
public enum QueryEngineEnum {
7+
HYPER("hyper"),
8+
TRINO("trino");
9+
10+
private final String value;
11+
12+
QueryEngineEnum(String value) {
13+
this.value = value;
14+
}
15+
16+
/**
17+
* get QueryEngineEnum value.
18+
*
19+
* @param value String value
20+
* @return QueryEngineEnum
21+
*/
22+
@JsonCreator
23+
public static QueryEngineEnum fromValue(String value) {
24+
for (QueryEngineEnum queryEngineEnum : QueryEngineEnum.values()) {
25+
if (queryEngineEnum.value.equals(value)) {
26+
return queryEngineEnum;
27+
}
28+
}
29+
30+
throw new IllegalArgumentException("Unexpected value for QueryEngineEnum: '" + value + "'");
31+
}
32+
33+
@Override
34+
@JsonValue
35+
public String toString() {
36+
return String.valueOf(this.value);
37+
}
38+
}
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
package com.salesforce.cdp.queryservice.model;
2+
3+
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
4+
import lombok.Data;
5+
6+
@Data
7+
@JsonIgnoreProperties(ignoreUnknown = true)
8+
public class QueryConfigResponse {
9+
String queryengine;
10+
}

src/main/java/com/salesforce/cdp/queryservice/util/Constants.java

Lines changed: 2 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ public class Constants {
2525
public static final String CDP_URL_V2 = "/api/v2";
2626
public static final String ANSI_SQL_URL = "/query";
2727
public static final String METADATA_URL = "/metadata";
28+
public static final String QUERY_CONFIG_URL = "/query-config";
2829
public static final String TOKEN_EXCHANGE_URL = "/services/a360/token";
2930
public static final String TOKEN_REVOKE_URL = "/services/oauth2/revoke";
3031
public static final String CORE_TOKEN_URL = "/services/oauth2/token";
@@ -34,18 +35,10 @@ public class Constants {
3435
public static final String DRIVER_NAME = "QueryService-jdbc";
3536
public static final String DRIVER_VERSION = "1.0";
3637

37-
// Common client id and secret information
38+
// Common server information
3839
public static final String PROD_SERVER_URL = ".salesforce.com";
39-
public static final String PROD_DEFAULT_CLIENT_ID = "3MVG9VeAQy5y3BQVJqaUbFmV5jd8imcck2K5idmrTTGocSu9qZZ6qkbuEkxECKVYwmzm3WgvxkujqsxZDcBpL";
40-
public static final String PROD_DEFAULT_CLIENT_SECRET = "1007FFFBA2B6B6B1EF21E2B03F4C4F692ADE10AB7DEF19E00D8AAF85EF6F6A12";
41-
4240
public static final String NA45_SERVER_URL = "na45.test1.pc-rnd.salesforce.com";
43-
public static final String NA45_DEFAULT_CLIENT_ID = "3MVG9XjhiDAzhaqaC4RR0yon8blsafwWlTnckUT8bEduWr0v9UpiQ2cJkmhNtFI1kVqpY8WpyE9JYkG.ZtgiE";
44-
public static final String NA45_DEFAULT_CLIENT_SECRET = "84642974065EDF90CA6F30FFEE23E2C16BDFA84D3083EF90E0AE905FA46131AD";
45-
4641
public static final String NA46_SERVER_URL = "na46.test1.pc-rnd.salesforce.com";
47-
public static final String NA46_DEFAULT_CLIENT_ID = "3MVG9sA57VMGPDfeS67yma6IPflHn83FRhxVpmnuzp7R8uS42JYshQ7gWgWR63CQRgKL9gY5AfitSme.01ib6";
48-
public static final String NA46_DEFAULT_CLIENT_SECRET = "BDAF015C3D2418008842CAE91B0C8DD2D672B41707FB11EA3CFC5A5392E31866";
4942

5043
//Audience constants for different environments
5144
public static final String PROD_SERVER_AUD = "login.salesforce.com";

src/main/java/com/salesforce/cdp/queryservice/util/Messages.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,4 +19,6 @@ public class Messages {
1919

2020
public static String RENEW_TOKEN = "Failed to Renew Token. Please retry";
2121

22+
public static String QUERY_CONFIG_ERROR = "Failed to get config from the server";
23+
2224
}

src/main/java/com/salesforce/cdp/queryservice/util/QueryExecutor.java

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,20 @@ public Response getMetadata() throws IOException, SQLException {
102102
return getResponse(request);
103103
}
104104

105+
public Response getQueryConfig() throws IOException, SQLException {
106+
log.info("Getting query config from CDP Query Service");
107+
Map<String, String> tokenWithTenantUrl = getTokenWithTenantUrl();
108+
String url = Constants.PROTOCOL + tokenWithTenantUrl.get(Constants.TENANT_URL)
109+
+ Constants.CDP_URL
110+
+ Constants.QUERY_CONFIG_URL;
111+
112+
Map<String, String> headers = createHeaders(tokenWithTenantUrl, false);
113+
headers.put(Constants.ENABLE_STREAM_FLOW, String.valueOf(this.connection.isEnableStreamFlow()));
114+
115+
Request request = HttpHelper.buildRequest(Constants.GET, url, null, headers);
116+
return getResponse(request);
117+
}
118+
105119
private Map<String, String> createHeaders(Map<String, String> tokenWithTenantUrl, boolean enableArrowStream) throws SQLException {
106120
Properties properties = connection.getClientInfo();
107121
Map<String, String> headers = new HashMap<>();
@@ -121,7 +135,7 @@ protected Response getResponse(Request request) throws IOException {
121135
// use queryClient to fetch metadata or to execute the query
122136
Response response = queryClient.newCall(request).execute();
123137
long endTime = System.currentTimeMillis();
124-
log.info("Total time taken to get response for url {} is {} ms", request.url(), endTime - startTime);
138+
log.info("Total time taken to get response for url {} is {} ms and traceid {}", request.url(), endTime - startTime, response.headers(Constants.TRACE_ID));
125139
return response;
126140
}
127141
}

src/main/java/com/salesforce/cdp/queryservice/util/QueryTokenExecutor.java

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -17,19 +17,16 @@
1717
package com.salesforce.cdp.queryservice.util;
1818

1919
import com.salesforce.cdp.queryservice.core.QueryServiceConnection;
20+
import com.salesforce.cdp.queryservice.enums.QueryEngineEnum;
2021
import com.salesforce.cdp.queryservice.interceptors.MetadataCacheInterceptor;
21-
import com.salesforce.cdp.queryservice.interceptors.RetryInterceptor;
2222
import com.salesforce.cdp.queryservice.model.Token;
2323
import com.salesforce.cdp.queryservice.util.internal.SFDefaultSocketFactoryWrapper;
2424
import lombok.extern.slf4j.Slf4j;
2525
import net.jodah.failsafe.Failsafe;
2626
import net.jodah.failsafe.FailsafeException;
2727
import net.jodah.failsafe.RetryPolicy;
2828
import okhttp3.OkHttpClient;
29-
import okhttp3.Request;
30-
import okhttp3.Response;
3129

32-
import java.io.IOException;
3330
import java.sql.SQLException;
3431
import java.util.Map;
3532
import java.util.Properties;
@@ -71,7 +68,7 @@ public QueryTokenExecutor(QueryServiceConnection connection, OkHttpClient client
7168
this.client = updateClientWithSocketFactory(client, connection.isSocksProxyDisabled());
7269

7370
// set TenantUrl in connection. This is mandatory in gRPC flow.
74-
if(connection.isEnableStreamFlow()) {
71+
if(QueryEngineEnum.HYPER == connection.getQueryEngineEnum()) {
7572
try {
7673
Map<String, String> token = getTokenWithTenantUrl();
7774
connection.setTenantUrl(token.get(Constants.TENANT_URL));

0 commit comments

Comments
 (0)