#8 SQL service refactoring

This commit is contained in:
serge-rider
2020-05-01 22:53:38 +03:00
parent 44ccbcf044
commit df8b2def9b
6 changed files with 250 additions and 133 deletions
@@ -17,11 +17,55 @@
package io.cloudbeaver.service.sql;
import io.cloudbeaver.DBWService;
import io.cloudbeaver.DBWebException;
import io.cloudbeaver.WebAction;
import io.cloudbeaver.model.WebAsyncTaskInfo;
import org.jkiss.code.NotNull;
import org.jkiss.code.Nullable;
import org.jkiss.dbeaver.DBException;
import org.jkiss.dbeaver.model.runtime.DBRProgressMonitor;
import org.jkiss.dbeaver.model.struct.DBSDataContainer;
import java.util.List;
import java.util.Map;
/**
* Web service API
* DBWServiceSQL
*/
public interface DBWServiceSQL extends DBWService {
@WebAction
WebSQLDialectInfo getDialectInfo(@NotNull WebSQLProcessor processor) throws DBWebException;
@WebAction
WebSQLCompletionProposal[] getCompletionProposals(@NotNull WebSQLContextInfo sqlContext, @NotNull String query, Integer position, Integer maxResults) throws DBWebException;
@WebAction
WebSQLContextInfo createContext(@NotNull WebSQLProcessor processor, String defaultCatalog, String defaultSchema);
@WebAction
void destroyContext(@NotNull WebSQLContextInfo sqlContext);
@WebAction
void setContextDefaults(@NotNull WebSQLContextInfo sqlContext, String catalogName, String schemaName) throws DBWebException;
@WebAction
@NotNull
WebSQLExecuteInfo executeQuery(@NotNull WebSQLContextInfo contextInfo, @NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBWebException;
@WebAction
Boolean closeResult(@NotNull WebSQLContextInfo sqlContext, @NotNull String resultId) throws DBWebException;
@WebAction
WebSQLExecuteInfo readDataFromContainer(@NotNull WebSQLContextInfo contextInfo, @NotNull String nodePath, @Nullable WebSQLDataFilter filter) throws DBException;
@WebAction
WebSQLExecuteInfo updateResultsData(
@NotNull WebSQLContextInfo contextInfo,
@NotNull String resultsId,
@NotNull List<Object> updateRow,
@NotNull Map<String, Object> updateValues) throws DBWebException;
@WebAction
WebAsyncTaskInfo asyncExecuteQuery(@NotNull WebSQLContextInfo contextInfo, @NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBException;
}
@@ -169,11 +169,6 @@ public class WebSQLContextInfo {
return CommonUtils.isEmpty(defaultSchema) ? null : defaultSchema;
}
@NotNull
public WebSQLExecuteInfo executeQuery(@NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBException {
return processor.executeQuery(this, sql, filter);
}
@NotNull
public WebSQLResultsInfo saveResult(@NotNull DBSDataContainer dataContainer, @NotNull DBDAttributeBinding[] attributes) {
WebSQLResultsInfo resultInfo = new WebSQLResultsInfo(
@@ -194,61 +189,14 @@ public class WebSQLContextInfo {
return resultsInfo;
}
@NotNull
public Boolean closeResult(@NotNull String resultId) {
public boolean closeResult(@NotNull String resultId) {
return resultInfoMap.remove(resultId) != null;
}
public @NotNull WebSQLCompletionProposal[] getCompletionProposals(@NotNull String query, Integer position, Integer maxResults) throws DBWebException {
try {
DBPDataSource dataSource = processor.getConnection().getDataSourceContainer().getDataSource();
Document document = new Document();
document.set(query);
SQLScriptElement activeQuery = new SQLQuery(dataSource, query);
SQLCompletionRequest request = new SQLCompletionRequest(
new WebSQLCompletionContext(this),
document,
position == null ? 0 : position,
activeQuery,
false
);
SQLCompletionAnalyzer analyzer = new SQLCompletionAnalyzer(request);
analyzer.runAnalyzer(processor.getWebSession().getProgressMonitor());
List<SQLCompletionProposalBase> proposals = analyzer.getProposals();
if (maxResults == null) maxResults = 200;
if (proposals.size() > maxResults) {
proposals = proposals.subList(0, maxResults);
}
WebSQLCompletionProposal[] result = new WebSQLCompletionProposal[proposals.size()];
for (int i = 0; i < proposals.size(); i++) {
result[i] = new WebSQLCompletionProposal(proposals.get(i));
}
return result;
} catch (DBException e) {
throw new DBWebException("Error processing SQL proposals", e);
}
}
///////////////////////////////////////////////////////
// Async model
@NotNull
public WebAsyncTaskInfo asyncExecuteQuery(@NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBException {
DBRRunnableWithResult<WebSQLExecuteInfo> runnable = new DBRRunnableWithResult<WebSQLExecuteInfo>() {
@Override
public void run(DBRProgressMonitor monitor) throws InvocationTargetException, InterruptedException {
try {
result = processor.processQuery(monitor, WebSQLContextInfo.this, sql, filter);
} catch (Throwable e) {
throw new InvocationTargetException(e);
}
}
};
return processor.getWebSession().createAndRunAsyncTask("SQL execute", runnable);
}
public void destroy() {
void dispose() {
resultInfoMap.clear();
}
@@ -19,30 +19,23 @@ package io.cloudbeaver.service.sql;
import io.cloudbeaver.DBWebException;
import io.cloudbeaver.WebAction;
import io.cloudbeaver.model.WebConnectionInfo;
import io.cloudbeaver.service.navigator.WebDatabaseObjectInfo;
import io.cloudbeaver.model.session.WebSession;
import io.cloudbeaver.service.navigator.WebStructContainers;
import org.jkiss.code.NotNull;
import org.jkiss.code.Nullable;
import org.jkiss.dbeaver.DBException;
import org.jkiss.dbeaver.Log;
import org.jkiss.dbeaver.model.DBPDataSource;
import org.jkiss.dbeaver.model.DBUtils;
import org.jkiss.dbeaver.model.data.*;
import org.jkiss.dbeaver.model.exec.*;
import org.jkiss.dbeaver.model.impl.data.DBDValueError;
import org.jkiss.dbeaver.model.impl.struct.ContextDefaultObjectsReader;
import org.jkiss.dbeaver.model.navigator.DBNDatabaseItem;
import org.jkiss.dbeaver.model.navigator.DBNNode;
import org.jkiss.dbeaver.model.runtime.DBRProgressMonitor;
import org.jkiss.dbeaver.model.sql.SQLDialect;
import org.jkiss.dbeaver.model.sql.SQLSyntaxManager;
import org.jkiss.dbeaver.model.sql.SQLUtils;
import org.jkiss.dbeaver.model.struct.DBSDataContainer;
import org.jkiss.dbeaver.model.struct.DBSDataManipulator;
import org.jkiss.dbeaver.model.struct.DBSEntity;
import org.jkiss.dbeaver.model.struct.DBSObject;
import org.jkiss.dbeaver.model.struct.rdb.DBSCatalog;
import org.jkiss.utils.CommonUtils;
import java.lang.reflect.InvocationTargetException;
@@ -68,7 +61,7 @@ public class WebSQLProcessor {
private AtomicInteger contextId = new AtomicInteger();
public WebSQLProcessor(@NotNull WebSession webSession, @NotNull WebConnectionInfo connection) {
WebSQLProcessor(@NotNull WebSession webSession, @NotNull WebConnectionInfo connection) {
this.webSession = webSession;
this.connection = connection;
@@ -76,6 +69,13 @@ public class WebSQLProcessor {
syntaxManager.init(connection.getDataSource());
}
void dispose() {
synchronized (contexts) {
contexts.forEach((s, context) -> context.dispose());
contexts.clear();
}
}
public WebConnectionInfo getConnection() {
return connection;
}
@@ -96,14 +96,6 @@ public class WebSQLProcessor {
return DBUtils.getDefaultContext(dataContainer, false);
}
@WebAction
@NotNull
public WebSQLDialectInfo getDialectInfo() throws DBWebException {
DBPDataSource dataSource = connection.getDataSourceContainer().getDataSource();
SQLDialect dialect = SQLUtils.getDialectFromDataSource(dataSource);
return new WebSQLDialectInfo(dataSource, dialect);
}
@NotNull
public WebSQLContextInfo createContext(String defaultCatalog, String defaultSchema) {
String contextId = String.valueOf(this.contextId.incrementAndGet());
@@ -128,24 +120,15 @@ public class WebSQLProcessor {
}
}
public void destroyContext(@NotNull String contextId) {
WebSQLContextInfo contextInfo;
public void destroyContext(@NotNull WebSQLContextInfo context) {
context.dispose();
synchronized (contexts) {
contextInfo = contexts.get(contextId);
}
if (contextInfo != null) {
contextInfo.destroy();
contexts.remove(context.getId());
}
}
@WebAction
@NotNull
public WebSQLExecuteInfo executeQuery(WebSQLContextInfo contextInfo, @NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBWebException {
return processQuery(webSession.getProgressMonitor(), contextInfo, sql, filter);
}
@NotNull
WebSQLExecuteInfo processQuery(DBRProgressMonitor monitor, WebSQLContextInfo contextInfo, @NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBWebException {
public WebSQLExecuteInfo processQuery(DBRProgressMonitor monitor, WebSQLContextInfo contextInfo, @NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBWebException {
if (filter == null) {
// Use default filter
filter = new WebSQLDataFilter();
@@ -193,28 +176,8 @@ public class WebSQLProcessor {
return executeInfo;
}
@WebAction
@NotNull
public WebSQLExecuteInfo readDataFromContainer(@NotNull WebSQLContextInfo contextInfo, @NotNull String containerPath, @Nullable WebSQLDataFilter filter) throws DBWebException {
try {
DBRProgressMonitor monitor = webSession.getProgressMonitor();
DBSDataContainer dataContainer = getDataContainerByNodePath(monitor, containerPath, DBSDataContainer.class);
if (filter == null) {
// Use default filter
filter = new WebSQLDataFilter();
}
return readDataFromContainer(contextInfo, monitor, dataContainer, filter);
} catch (DBException e) {
if (e instanceof DBWebException) throw (DBWebException) e;
throw new DBWebException("Error reading data from '" + containerPath + "'", e);
}
}
@NotNull
private WebSQLExecuteInfo readDataFromContainer(@NotNull WebSQLContextInfo contextInfo, @NotNull DBRProgressMonitor monitor, @NotNull DBSDataContainer dataContainer, @NotNull WebSQLDataFilter filter) throws DBException {
public WebSQLExecuteInfo readDataFromContainer(@NotNull WebSQLContextInfo contextInfo, @NotNull DBRProgressMonitor monitor, @NotNull DBSDataContainer dataContainer, @NotNull WebSQLDataFilter filter) throws DBException {
WebSQLExecuteInfo executeInfo = new WebSQLExecuteInfo();
@@ -41,10 +41,11 @@ public class WebServiceBindingSQL extends WebServiceBindingBase<DBWServiceSQL> {
public void bindWiring(DBWBindingContext model) throws DBWebException {
model.getQueryType()
.dataFetcher("sqlDialectInfo", env ->
getSQLProcessor(model, env).getDialectInfo()
getService(env).getDialectInfo(getSQLProcessor(env))
)
.dataFetcher("sqlCompletionProposals", env ->
getSQLContext(model, env).getCompletionProposals(
getService(env).getCompletionProposals(
getSQLContext(env),
env.getArgument("query"),
env.getArgument("position"),
env.getArgument("maxResults")
@@ -52,46 +53,62 @@ public class WebServiceBindingSQL extends WebServiceBindingBase<DBWServiceSQL> {
);
model.getMutationType()
.dataFetcher("sqlContextCreate", env -> getSQLProcessor(model, env).createContext(env.getArgument("defaultCatalog"), env.getArgument("defaultSchema")))
.dataFetcher("sqlContextDestroy", env -> { getSQLContext(model, env).destroy(); return true; } )
.dataFetcher("sqlContextSetDefaults", env -> { getSQLContext(model, env).setDefaults(env.getArgument("defaultCatalog"), env.getArgument("defaultSchema")); return true; })
.dataFetcher("sqlContextCreate", env -> getService(env).createContext(
getSQLProcessor(env),
env.getArgument("defaultCatalog"),
env.getArgument("defaultSchema")))
.dataFetcher("sqlContextDestroy", env -> { getService(env).destroyContext(getSQLContext(env)); return true; } )
.dataFetcher("sqlContextSetDefaults", env -> {
getService(env).setContextDefaults(
getSQLContext(env),
env.getArgument("defaultCatalog"),
env.getArgument("defaultSchema"));
return true;
})
.dataFetcher("sqlExecuteQuery", env ->
getSQLContext(model, env).executeQuery(
env.getArgument("sql"), getDataFilter(env)
getService(env).executeQuery(
getSQLContext(env),
env.getArgument("sql"),
getDataFilter(env)
))
.dataFetcher("sqlResultClose", env ->
getSQLContext(model, env).closeResult(env.getArgument("resultId")))
getService(env).closeResult(
getSQLContext(env),
env.getArgument("resultId")))
.dataFetcher("readDataFromContainer", env ->
getSQLProcessor(model, env).readDataFromContainer(
getSQLContext(model, env),
env.getArgument("containerNodePath"), getDataFilter(env)
getService(env).readDataFromContainer(
getSQLContext(env),
env.getArgument("containerNodePath"),
getDataFilter(env)
))
.dataFetcher("updateResultsData", env ->
getSQLProcessor(model, env).updateResultsData(
getSQLContext(model, env),
getService(env).updateResultsData(
getSQLContext(env),
env.getArgument("resultsId"),
env.getArgument("updateRow"),
env.getArgument("updateValues")
))
.dataFetcher("asyncSqlExecuteQuery", env ->
getSQLContext(model, env).asyncExecuteQuery(
env.getArgument("sql"), getDataFilter(env)
getService(env).asyncExecuteQuery(
getSQLContext(env),
env.getArgument("sql"),
getDataFilter(env)
));
}
public static WebSQLConfiguration getSQLConfiguration(WebSession webSession) {
return webSession.getAttribute("sqlConfiguration", cfg -> new WebSQLConfiguration(), cfg -> null);
return webSession.getAttribute("sqlConfiguration", cfg -> new WebSQLConfiguration(), WebSQLConfiguration::dispose);
}
public static WebSQLProcessor getSQLProcessor(DBWBindingContext model, DataFetchingEnvironment env) throws DBWebException {
public static WebSQLProcessor getSQLProcessor(DataFetchingEnvironment env) throws DBWebException {
WebConnectionInfo connectionInfo = getWebConnection(env);
return getSQLConfiguration(connectionInfo.getSession()).getSQLProcessor(connectionInfo);
}
public static WebSQLContextInfo getSQLContext(DBWBindingContext model, DataFetchingEnvironment env) throws DBWebException {
WebSQLProcessor processor = getSQLProcessor(model, env);
public static WebSQLContextInfo getSQLContext(DataFetchingEnvironment env) throws DBWebException {
WebSQLProcessor processor = getSQLProcessor(env);
String contextId = env.getArgument("contextId");
WebSQLContextInfo context = processor.getContext(contextId);
if (context == null) {
@@ -111,12 +128,22 @@ public class WebServiceBindingSQL extends WebServiceBindingBase<DBWServiceSQL> {
throw new DBWebException("Error connecting to database", e);
}
}
WebSQLProcessor processor = processors.get(connectionInfo);
if (processor == null) {
processor = new WebSQLProcessor(connectionInfo.getSession(), connectionInfo);
processors.put(connectionInfo, processor);
synchronized (processors) {
WebSQLProcessor processor = processors.get(connectionInfo);
if (processor == null) {
processor = new WebSQLProcessor(connectionInfo.getSession(), connectionInfo);
processors.put(connectionInfo, processor);
}
return processor;
}
return processor;
}
public WebSQLConfiguration dispose() {
synchronized (processors) {
processors.forEach((connectionInfo, processor) -> processor.dispose());
processors.clear();
}
return this;
}
}
@@ -17,11 +17,146 @@
package io.cloudbeaver.service.sql.impl;
import io.cloudbeaver.service.sql.DBWServiceSQL;
import io.cloudbeaver.DBWebException;
import io.cloudbeaver.WebAction;
import io.cloudbeaver.model.WebAsyncTaskInfo;
import io.cloudbeaver.service.sql.*;
import org.eclipse.jface.text.Document;
import org.jkiss.code.NotNull;
import org.jkiss.code.Nullable;
import org.jkiss.dbeaver.DBException;
import org.jkiss.dbeaver.model.DBPDataSource;
import org.jkiss.dbeaver.model.runtime.DBRProgressMonitor;
import org.jkiss.dbeaver.model.runtime.DBRRunnableWithResult;
import org.jkiss.dbeaver.model.sql.SQLDialect;
import org.jkiss.dbeaver.model.sql.SQLQuery;
import org.jkiss.dbeaver.model.sql.SQLScriptElement;
import org.jkiss.dbeaver.model.sql.SQLUtils;
import org.jkiss.dbeaver.model.sql.completion.SQLCompletionAnalyzer;
import org.jkiss.dbeaver.model.sql.completion.SQLCompletionProposalBase;
import org.jkiss.dbeaver.model.sql.completion.SQLCompletionRequest;
import org.jkiss.dbeaver.model.struct.DBSDataContainer;
import java.lang.reflect.InvocationTargetException;
import java.util.List;
import java.util.Map;
/**
* Web service implementation
*/
public class WebServiceSQL implements DBWServiceSQL {
@Override
@NotNull
public WebSQLDialectInfo getDialectInfo(@NotNull WebSQLProcessor processor) throws DBWebException {
DBPDataSource dataSource = processor.getConnection().getDataSourceContainer().getDataSource();
SQLDialect dialect = SQLUtils.getDialectFromDataSource(dataSource);
return new WebSQLDialectInfo(dataSource, dialect);
}
@NotNull
public WebSQLCompletionProposal[] getCompletionProposals(@NotNull WebSQLContextInfo sqlContext, @NotNull String query, Integer position, Integer maxResults) throws DBWebException {
try {
DBPDataSource dataSource = sqlContext.getProcessor().getConnection().getDataSourceContainer().getDataSource();
Document document = new Document();
document.set(query);
SQLScriptElement activeQuery = new SQLQuery(dataSource, query);
SQLCompletionRequest request = new SQLCompletionRequest(
new WebSQLCompletionContext(sqlContext),
document,
position == null ? 0 : position,
activeQuery,
false
);
SQLCompletionAnalyzer analyzer = new SQLCompletionAnalyzer(request);
analyzer.runAnalyzer(sqlContext.getProcessor().getWebSession().getProgressMonitor());
List<SQLCompletionProposalBase> proposals = analyzer.getProposals();
if (maxResults == null) maxResults = 200;
if (proposals.size() > maxResults) {
proposals = proposals.subList(0, maxResults);
}
WebSQLCompletionProposal[] result = new WebSQLCompletionProposal[proposals.size()];
for (int i = 0; i < proposals.size(); i++) {
result[i] = new WebSQLCompletionProposal(proposals.get(i));
}
return result;
} catch (DBException e) {
throw new DBWebException("Error processing SQL proposals", e);
}
}
@Override
public WebSQLContextInfo createContext(@NotNull WebSQLProcessor processor, String defaultCatalog, String defaultSchema) {
return processor.createContext(defaultCatalog, defaultSchema);
}
@Override
public void destroyContext(@NotNull WebSQLContextInfo sqlContext) {
sqlContext.getProcessor().destroyContext(sqlContext);
}
@Override
public void setContextDefaults(@NotNull WebSQLContextInfo sqlContext, String catalogName, String schemaName) throws DBWebException {
sqlContext.setDefaults(catalogName, schemaName);
}
@WebAction
@NotNull
public WebSQLExecuteInfo executeQuery(@NotNull WebSQLContextInfo sqlContext, @NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBWebException {
return sqlContext.getProcessor().processQuery(
sqlContext.getProcessor().getWebSession().getProgressMonitor(),
sqlContext,
sql,
filter);
}
@Override
public Boolean closeResult(@NotNull WebSQLContextInfo sqlContext, @NotNull String resultId) throws DBWebException {
if (!sqlContext.closeResult(resultId)) {
throw new DBWebException("Invalid result ID " + resultId);
}
return true;
}
@Override
public WebSQLExecuteInfo readDataFromContainer(@NotNull WebSQLContextInfo contextInfo, @NotNull String containerPath, @Nullable WebSQLDataFilter filter) throws DBException {
try {
DBRProgressMonitor monitor = contextInfo.getProcessor().getWebSession().getProgressMonitor();
DBSDataContainer dataContainer = contextInfo.getProcessor().getDataContainerByNodePath(monitor, containerPath, DBSDataContainer.class);
if (filter == null) {
// Use default filter
filter = new WebSQLDataFilter();
}
return contextInfo.getProcessor().readDataFromContainer(contextInfo, monitor, dataContainer, filter);
} catch (DBException e) {
if (e instanceof DBWebException) throw (DBWebException) e;
throw new DBWebException("Error reading data from '" + containerPath + "'", e);
}
}
@Override
public WebSQLExecuteInfo updateResultsData(@NotNull WebSQLContextInfo contextInfo, @NotNull String resultsId, @NotNull List<Object> updateRow, @NotNull Map<String, Object> updateValues) throws DBWebException {
return contextInfo.getProcessor().updateResultsData(contextInfo, resultsId, updateRow, updateValues);
}
@NotNull
public WebAsyncTaskInfo asyncExecuteQuery(@NotNull WebSQLContextInfo contextInfo, @NotNull String sql, @Nullable WebSQLDataFilter filter) throws DBException {
DBRRunnableWithResult<WebSQLExecuteInfo> runnable = new DBRRunnableWithResult<WebSQLExecuteInfo>() {
@Override
public void run(DBRProgressMonitor monitor) throws InvocationTargetException, InterruptedException {
try {
result = contextInfo.getProcessor().processQuery(monitor, contextInfo, sql, filter);
} catch (Throwable e) {
throw new InvocationTargetException(e);
}
}
};
return contextInfo.getProcessor().getWebSession().createAndRunAsyncTask("SQL execute", runnable);
}
}
@@ -43,17 +43,17 @@ public class WebServiceBindingDataTransfer extends WebServiceBindingBase<DBWServ
.dataFetcher("dataTransferAvailableStreamProcessors",
env -> getService(env).getAvailableStreamProcessors(getWebSession(env)))
.dataFetcher("dataTransferExportDataFromContainer", env -> getService(env).dataTransferExportDataFromContainer(
WebServiceBindingSQL.getSQLProcessor(model, env),
WebServiceBindingSQL.getSQLProcessor(env),
env.getArgument("containerNodePath"),
new WebDataTransferParameters(env.getArgument("parameters"))
))
.dataFetcher("dataTransferExportDataFromResults", env -> getService(env).dataTransferExportDataFromResults(
WebServiceBindingSQL.getSQLContext(model, env),
WebServiceBindingSQL.getSQLContext(env),
env.getArgument("resultsId"),
new WebDataTransferParameters(env.getArgument("parameters"))
))
.dataFetcher("dataTransferRemoveDataFile", env -> getService(env).dataTransferRemoveDataFile(
WebServiceBindingSQL.getSQLProcessor(model, env),
WebServiceBindingSQL.getSQLProcessor(env),
env.getArgument("dataFileId")
))
;