001package gudusoft.gsqlparser.dlineage.dataflow.model;
002
003import gudusoft.gsqlparser.EDbVendor;
004import gudusoft.gsqlparser.ESetOperatorType;
005import gudusoft.gsqlparser.TCustomSqlStatement;
006import gudusoft.gsqlparser.TSourceToken;
007import gudusoft.gsqlparser.dlineage.util.DlineageUtil;
008import gudusoft.gsqlparser.dlineage.util.Pair3;
009import gudusoft.gsqlparser.nodes.TCTE;
010import gudusoft.gsqlparser.nodes.TFunctionCall;
011import gudusoft.gsqlparser.nodes.TJoinList;
012import gudusoft.gsqlparser.pp.utils.SourceTokenSearcher;
013import gudusoft.gsqlparser.sqlenv.ESQLDataObjectType;
014import gudusoft.gsqlparser.sqlenv.TSQLEnv;
015import gudusoft.gsqlparser.stmt.*;
016import gudusoft.gsqlparser.stmt.hive.THiveLoad;
017import gudusoft.gsqlparser.stmt.mssql.TCreateExternalDataSourceStmt;
018import gudusoft.gsqlparser.stmt.mssql.TMssqlCreateFunction;
019import gudusoft.gsqlparser.stmt.mssql.TMssqlDeclare;
020import gudusoft.gsqlparser.stmt.mssql.TMssqlExecute;
021import gudusoft.gsqlparser.stmt.snowflake.TCreateStageStmt;
022import gudusoft.gsqlparser.stmt.snowflake.TCreateStreamStmt;
023import gudusoft.gsqlparser.stmt.snowflake.TSnowflakeCopyIntoStmt;
024import gudusoft.gsqlparser.stmt.teradata.TTeradataCreateProcedure;
025import gudusoft.gsqlparser.util.Logger;
026import gudusoft.gsqlparser.util.LoggerFactory;
027import gudusoft.gsqlparser.util.SQLUtil;
028
029import java.io.ByteArrayInputStream;
030import java.io.IOException;
031import java.util.ArrayDeque;
032import java.util.ArrayList;
033import java.util.Deque;
034import java.util.List;
035import java.util.Properties;
036
037public class Process {
038    private static final Logger logger = LoggerFactory.getLogger(Process.class);
039    private long id;
040    protected String server;
041    protected String schema;
042    protected String database;
043    private String name;
044    private String procedureName;
045    private Long procedureId;
046    private String queryHashId;
047    private String type;
048    private String customType;
049    private TCustomSqlStatement gspObject;
050    private Pair3<Long, Long, String> startPosition;
051    private Pair3<Long, Long, String> endPosition;
052    private List<Object> targetColumns = new ArrayList<Object>();
053    private List<Transform> transforms = new ArrayList<Transform>();
054
055    public Process(TCustomSqlStatement gspObject) {
056
057        if (gspObject == null) {
058            throw new IllegalArgumentException("Process arguments can't be null.");
059        }
060
061        id = ++ModelBindingManager.get().TABLE_COLUMN_ID;
062
063        this.gspObject = gspObject;
064
065        TSourceToken startToken = gspObject.getStartToken();
066        TSourceToken endToken = gspObject.getEndToken();
067        if (startToken != null) {
068            this.startPosition = new Pair3<Long, Long, String>(startToken.lineNo, startToken.columnNo,
069                    ModelBindingManager.getGlobalHash());
070        }
071
072        if (endToken != null) {
073            this.endPosition = new Pair3<Long, Long, String>(endToken.lineNo,
074                    endToken.columnNo + SQLUtil.endTrim(endToken.getAstext()).length(), ModelBindingManager.getGlobalHash());
075        }
076
077        this.schema = ModelBindingManager.getGlobalSchema();
078        this.database = ModelBindingManager.getGlobalDatabase();
079        this.queryHashId = SQLUtil.stringToMD5(gspObject.asCanonical().replaceAll("\r?\n", "\r\n"));
080
081        EDbVendor vendor = ModelBindingManager.getGlobalOption().getVendor();
082        boolean supportCatalog = TSQLEnv.supportCatalog(vendor);
083        boolean supportSchema = TSQLEnv.supportSchema(vendor);
084
085        fillSchemaInfo();
086
087        if (!supportCatalog) {
088            this.database = null;
089        } else if (this.database == null && !TSQLEnv.DEFAULT_DB_NAME.equals(getDefaultDatabase())) {
090            this.database = getDefaultDatabase();
091        }
092
093        if (!supportSchema) {
094            this.schema = null;
095        } else if (this.schema == null && !TSQLEnv.DEFAULT_SCHEMA_NAME.equals(getDefaultSchema())) {
096            this.schema = getDefaultSchema();
097        }
098
099        String procedureParent = getProcedureParentName(gspObject);
100        if (procedureParent != null) {
101            procedureName = procedureParent;
102            Procedure procedure = ModelBindingManager.get()
103                    .getProcedureByName(DlineageUtil.getIdentifierNormalTableName(
104                            DlineageUtil.getProcedureNameWithArgs(getParentProcedure(gspObject))));
105            if (procedure == null) {
106                procedure = ModelBindingManager.get().getProcedureByName(DlineageUtil.getIdentifierNormalTableName(
107                        DlineageUtil.getProcedureNameWithArgNum(getParentProcedure(gspObject))));
108            }
109            if (procedure != null) {
110                this.procedureId = procedure.getId();
111            }
112        } else {
113            procedureName = "batchQueries";
114        }
115
116        TCustomSqlStatement stmt = DlineageUtil.getTopStmt(ModelBindingManager.getGlobalStmtStack().peek());
117        String sqlComment = null;
118        try {
119            sqlComment = stmt.getCommentBeforeNode();
120        } catch (Exception e) {
121        }
122        if (!SQLUtil.isEmpty(sqlComment) && sqlComment.indexOf("process") != -1) {
123            Properties properties = new Properties();
124            try {
125                properties.load(
126                        new ByteArrayInputStream(sqlComment.replace("--", "").trim().replace(",", "\n").getBytes()));
127                if (properties.containsKey("process_label")) {
128                    this.customType = properties.getProperty("process_label");
129                } else if (properties.containsKey("process")) {
130                    this.customType = properties.getProperty("process");
131                }
132            } catch (IOException e) {
133                logger.error("load sql comment properties failed.", e);
134            }
135        }
136
137        if (gspObject instanceof TCreateTableSqlStatement) {
138            this.type = "Create Table";
139        } else if (gspObject instanceof TCreateViewSqlStatement) {
140            this.type = "Create View";
141        } else if (gspObject instanceof TUpdateSqlStatement) {
142            this.type = "Update";
143        } else if (gspObject instanceof TMergeSqlStatement) {
144            this.type = "Merge";
145        } else if (gspObject instanceof TInsertSqlStatement) {
146            this.type = "Insert";
147        } else if (gspObject instanceof TSelectSqlStatement
148                && ((TSelectSqlStatement) gspObject).getIntoClause() != null) {
149            this.type = "Select Into";
150        } else if (gspObject instanceof TCreateTriggerStmt) {
151            this.type = "Create Trigger";
152        } else if (gspObject instanceof TCreateFunctionStmt) {
153            this.type = "Create Function";
154        } else if (gspObject instanceof TMssqlCreateFunction) {
155            this.type = "Create Function";
156        } else if (gspObject instanceof TCreateProcedureStmt) {
157            this.type = "Create Procedure";
158        } else if (gspObject instanceof THiveLoad) {
159            this.type = "Hive Load";
160        } else if (gspObject instanceof TSnowflakeCopyIntoStmt) {
161            this.type = "Copy Into";
162        } else if (gspObject instanceof TAlterTableStatement) {
163            this.type = "Alter Table";
164        } else if (gspObject instanceof TRenameStmt) {
165            this.type = "Rename Table";
166        } else if (gspObject instanceof TMssqlDeclare) {
167            this.type = "Mssql Declare";
168        } else if (gspObject instanceof TCreateStageStmt) {
169            this.type = "Create Stage";
170        } else if (gspObject instanceof TCreateSynonymStmt) {
171            this.type = "Create Synonym";
172        } else if (gspObject instanceof TCreateExternalDataSourceStmt) {
173            this.type = "Create Datasource";
174        } else if (gspObject instanceof TCreateStreamStmt) {
175            this.type = "Create Stream";
176        } else if (gspObject instanceof TCreateDatabaseSqlStatement) {
177            this.type = "Create Database";
178        } else if (gspObject instanceof TCreateSchemaSqlStatement) {
179            this.type = "Create Schema";
180        } else if (gspObject instanceof TMssqlExecute && ((TMssqlExecute)gspObject).getModuleName().toString().equalsIgnoreCase("sp_rename")) {
181            this.type = "Rename Table";
182        } else if (gspObject instanceof TCallStatement){
183            this.type = "Function Call";
184        } else {
185            this.type = gspObject.getClass().getSimpleName();
186        }
187
188        if (this.server == null && !TSQLEnv.DEFAULT_SERVER_NAME.equals(getDefaultServer())) {
189            this.server = getDefaultServer();
190        }
191
192        if (stmt.getStatements() != null) {
193            for (int i = 0; i < stmt.getStatements().size(); i++) {
194                TCustomSqlStatement subquery = stmt.getStatements().get(i);
195                if (subquery instanceof TSelectSqlStatement) {
196                    TSelectSqlStatement select = (TSelectSqlStatement) subquery;
197                    appendTransform(select, ESetOperatorType.none);
198                }
199            }
200        }
201    }
202
203    public Process(TFunctionCall gspObject) {
204
205        if (gspObject == null) {
206            throw new IllegalArgumentException("Process arguments can't be null.");
207        }
208
209        id = ++ModelBindingManager.get().TABLE_COLUMN_ID;
210
211
212        TSourceToken startToken = gspObject.getStartToken();
213        TSourceToken endToken = gspObject.getEndToken();
214        if (startToken != null) {
215            this.startPosition = new Pair3<Long, Long, String>(startToken.lineNo, startToken.columnNo,
216                    ModelBindingManager.getGlobalHash());
217        }
218
219        if (endToken != null) {
220            this.endPosition = new Pair3<Long, Long, String>(endToken.lineNo,
221                    endToken.columnNo + SQLUtil.endTrim(endToken.getAstext()).length(), ModelBindingManager.getGlobalHash());
222        }
223
224        TCustomSqlStatement stmt = DlineageUtil.getTopStmt(ModelBindingManager.getGlobalStmtStack().peek());
225        this.gspObject = stmt;
226
227        this.schema = ModelBindingManager.getGlobalSchema();
228        this.database = ModelBindingManager.getGlobalDatabase();
229
230        if (stmt != null) {
231            this.queryHashId = SQLUtil.stringToMD5(stmt.asCanonical().replaceAll("\r?\n", "\r\n"));
232        }
233
234        EDbVendor vendor = ModelBindingManager.getGlobalOption().getVendor();
235        boolean supportCatalog = TSQLEnv.supportCatalog(vendor);
236        boolean supportSchema = TSQLEnv.supportSchema(vendor);
237
238        fillSchemaInfo();
239
240        if (!supportCatalog) {
241            this.database = null;
242        } else if (this.database == null && !TSQLEnv.DEFAULT_DB_NAME.equals(getDefaultDatabase())) {
243            this.database = getDefaultDatabase();
244        }
245
246        if (!supportSchema) {
247            this.schema = null;
248        } else if (this.schema == null && !TSQLEnv.DEFAULT_SCHEMA_NAME.equals(getDefaultSchema())) {
249            this.schema = getDefaultSchema();
250        }
251
252        String procedureParent = getProcedureParentName(stmt);
253        if (procedureParent != null) {
254            procedureName = procedureParent;
255            Procedure procedure = ModelBindingManager.get()
256                    .getProcedureByName(DlineageUtil.getIdentifierNormalTableName(
257                            DlineageUtil.getProcedureNameWithArgs(getParentProcedure(stmt))));
258            if (procedure == null) {
259                procedure = ModelBindingManager.get().getProcedureByName(DlineageUtil.getIdentifierNormalTableName(
260                        DlineageUtil.getProcedureNameWithArgNum(getParentProcedure(stmt))));
261            }
262            if (procedure != null) {
263                this.procedureId = procedure.getId();
264            }
265        } else {
266            procedureName = "batchQueries";
267        }
268
269        String sqlComment = null;
270        try {
271            sqlComment = stmt.getCommentBeforeNode();
272        } catch (Exception e) {
273        }
274        if (!SQLUtil.isEmpty(sqlComment) && sqlComment.indexOf("process") != -1) {
275            Properties properties = new Properties();
276            try {
277                properties.load(
278                        new ByteArrayInputStream(sqlComment.replace("--", "").trim().replace(",", "\n").getBytes()));
279                if (properties.containsKey("process_label")) {
280                    this.customType = properties.getProperty("process_label");
281                } else if (properties.containsKey("process")) {
282                    this.customType = properties.getProperty("process");
283                }
284            } catch (IOException e) {
285                logger.error("load sql comment properties failed.", e);
286            }
287        }
288
289        this.type = "Function Call";
290
291        if (this.server == null && !TSQLEnv.DEFAULT_SERVER_NAME.equals(getDefaultServer())) {
292            this.server = getDefaultServer();
293        }
294    }
295
296    private void appendTransform(TSelectSqlStatement select, ESetOperatorType operatorType) {
297        // Iterative traversal of the UNION tree to avoid StackOverflow with deeply nested unions.
298        Deque<Object[]> stack = new ArrayDeque<>();
299        stack.push(new Object[]{select, operatorType});
300
301        while (!stack.isEmpty()) {
302            Object[] entry = stack.pop();
303            TSelectSqlStatement current = (TSelectSqlStatement) entry[0];
304            ESetOperatorType currentOp = (ESetOperatorType) entry[1];
305
306            if (currentOp != ESetOperatorType.none) {
307                Transform transform = new Transform();
308                transform.setType(currentOp.name());
309                transform.setCodeString(current.toString());
310                transform.setStartToken(current.getStartToken());
311                transform.setEndToken(current.getEndToken());
312                transforms.add(transform);
313            } else {
314                if (current.getCteList() != null && current.getCteList().size() > 0) {
315                    for (int i = 0; i < current.getCteList().size(); i++) {
316                        TCTE cte = current.getCteList().getCTE(i);
317                        Transform transform = new Transform();
318                        transform.setType(Transform.CTE);
319                        transform.setCodeString(cte.toString());
320                        transform.setStartToken(cte.getStartToken());
321                        transform.setEndToken(cte.getEndToken());
322                        transforms.add(transform);
323                    }
324                }
325                if (current.getSetOperatorType() != ESetOperatorType.none) {
326                    // Push right first so left is processed first (stack is LIFO)
327                    if (current.getRightStmt() != null) {
328                        ESetOperatorType rightOp = current.getRightStmt().getSetOperatorType() != ESetOperatorType.none
329                                ? ESetOperatorType.none : current.getSetOperatorType();
330                        stack.push(new Object[]{current.getRightStmt(), rightOp});
331                    }
332                    if (current.getLeftStmt() != null) {
333                        ESetOperatorType leftOp = current.getLeftStmt().getSetOperatorType() != ESetOperatorType.none
334                                ? ESetOperatorType.none : current.getSetOperatorType();
335                        stack.push(new Object[]{current.getLeftStmt(), leftOp});
336                    }
337                } else {
338                    if (current.getJoins() != null && current.getJoins().size() > 0) {
339                        TJoinList joins = current.getJoins();
340                        TSourceToken startToken = joins.getStartToken();
341                        TSourceToken fromToken = SourceTokenSearcher.backforwardSearch(startToken, 10, "from");
342                        if (fromToken == null) {
343                            if (startToken.getAstext().equalsIgnoreCase("from")) {
344                                fromToken = startToken;
345                            } else {
346                                continue;
347                            }
348                        }
349                        StringBuilder builder = new StringBuilder();
350                        if (joins.getEndToken().posinlist < fromToken.container.size()
351                                && joins.getEndToken().posinlist >= fromToken.posinlist) {
352                            for (int j = fromToken.posinlist; j <= joins.getEndToken().posinlist; j++) {
353                                builder.append(fromToken.container.get(j));
354                            }
355                        } else {
356                            System.err.println("Handle statement transform error, statemet is:");
357                            System.err.println(current.toString());
358                        }
359                        Transform transform = new Transform();
360                        transform.setType(Transform.FROM);
361                        transform.setCodeString(builder.toString());
362                        transform.setStartToken(fromToken);
363                        transform.setEndToken(joins.getEndToken());
364                        transforms.add(transform);
365                    }
366                }
367            }
368        }
369    }
370
371    private TSelectSqlStatement getFirstSubquery(TSelectSqlStatement select) {
372        // Iterative left-chain walk to avoid StackOverflow with deeply nested unions.
373        TSelectSqlStatement current = select;
374        while (current.getSetOperatorType() != ESetOperatorType.none) {
375            current = current.getLeftStmt();
376        }
377        return current;
378    }
379
380    private TSelectSqlStatement getLastSubquery(TSelectSqlStatement select) {
381        // Iterative right-chain walk to avoid StackOverflow with deeply nested unions.
382        TSelectSqlStatement current = select;
383        while (current.getSetOperatorType() != ESetOperatorType.none) {
384            current = current.getRightStmt();
385        }
386        return current;
387    }
388
389    private void fillSchemaInfo() {
390        TCustomSqlStatement stmt = DlineageUtil.getTopStmt(ModelBindingManager.getGlobalStmtStack().peek());
391        String sqlComment = null;
392        try {
393            sqlComment = stmt.getCommentBeforeNode();
394        } catch (Exception e) {
395        }
396        if (!SQLUtil.isEmpty(sqlComment) && (sqlComment.indexOf("db") != -1 || sqlComment.indexOf("schema") != -1)) {
397            Properties properties = new Properties();
398            try {
399                properties.load(
400                        new ByteArrayInputStream(sqlComment.replace("--", "").trim().replace(",", "\n").getBytes()));
401                if (SQLUtil.isEmpty(this.server) && properties.containsKey("db-instance")) {
402                    this.server = properties.getProperty("db-instance");
403                }
404                if (SQLUtil.isEmpty(this.database) && properties.containsKey("db")) {
405                    this.database = properties.getProperty("db");
406                    if (this.database.indexOf(".") != -1) {
407                        this.database = SQLUtil.quoteDottedName(ModelBindingManager.getGlobalOption().getVendor(), ESQLDataObjectType.dotCatalog, this.database);
408                    }
409                }
410                if (SQLUtil.isEmpty(this.schema) && properties.containsKey("schema")) {
411                    this.schema = properties.getProperty("schema");
412                    if (this.schema.indexOf(".") != -1) {
413                        this.schema = SQLUtil.quoteDottedName(ModelBindingManager.getGlobalOption().getVendor(), ESQLDataObjectType.dotSchema, this.schema);
414                    }
415                }
416            } catch (IOException e) {
417                logger.error("load sql comment properties failed.", e);
418            }
419        }
420    }
421
422    protected String getDefaultServer() {
423        String defaultServer = null;
424        if (ModelBindingManager.getGlobalSQLEnv() != null) {
425            defaultServer = ModelBindingManager.getGlobalSQLEnv().getDefaultServerName();
426        }
427        if (!SQLUtil.isEmpty(defaultServer))
428            return defaultServer;
429        return TSQLEnv.DEFAULT_SERVER_NAME;
430    }
431
432    protected String getDefaultSchema() {
433        return ModelBindingManager.getEffectiveDefaultSchema();
434    }
435
436    protected String getDefaultDatabase() {
437        return ModelBindingManager.getEffectiveDefaultDatabase();
438    }
439
440    private String getProcedureParentName(TCustomSqlStatement stmt) {
441        if (stmt instanceof TStoredProcedureSqlStatement) {
442                        if (((TStoredProcedureSqlStatement) stmt).getStoredProcedureName() != null) {
443                                return ((TStoredProcedureSqlStatement) stmt).getStoredProcedureName().toString();
444                        }
445                }
446
447        if (stmt == null) {
448                return null;
449        }
450        
451        if (stmt instanceof TTeradataCreateProcedure) {
452            if (((TTeradataCreateProcedure) stmt).getProcedureName() != null) {
453                return ((TTeradataCreateProcedure) stmt).getProcedureName().toString();
454            }
455        }
456        
457        if (stmt instanceof TStoredProcedureSqlStatement) {
458            if (((TStoredProcedureSqlStatement) stmt).getStoredProcedureName() != null) {
459                return ((TStoredProcedureSqlStatement) stmt).getStoredProcedureName().toString();
460            }
461
462            if (stmt instanceof TCommonBlock) {
463                if (((TCommonBlock) stmt).getBlockBody().getParentObjectName() instanceof TStoredProcedureSqlStatement) {
464                    stmt = (TStoredProcedureSqlStatement) ((TCommonBlock) stmt).getBlockBody().getParentObjectName();
465                }
466                else {
467                        return getProcedureParentName(stmt.getParentStmt());
468                }
469            }
470            
471            return null;
472        }
473        
474                stmt = stmt.getParentStmt();
475        if (stmt == null) {
476            if(ModelBindingManager.getGlobalProcedure()!=null) {
477                return getProcedureParentName(ModelBindingManager.getGlobalProcedure());
478            }
479            else {
480                return null;
481            }
482        }
483        
484
485        if (stmt instanceof TCommonBlock) {
486            if (((TCommonBlock) stmt).getBlockBody().getParentObjectName() instanceof TStoredProcedureSqlStatement) {
487                stmt = (TStoredProcedureSqlStatement) ((TCommonBlock) stmt).getBlockBody().getParentObjectName();
488            }
489        }
490
491        return getProcedureParentName(stmt);
492    }
493
494    private TStoredProcedureSqlStatement getParentProcedure(TCustomSqlStatement stmt) {
495        if (stmt == null) {
496                return null;
497        }
498        
499        if (stmt instanceof TStoredProcedureSqlStatement) {
500            if (((TStoredProcedureSqlStatement) stmt).getStoredProcedureName() != null) {
501                return ((TStoredProcedureSqlStatement) stmt);
502            }
503            
504            if (stmt instanceof TCommonBlock) {
505                if (((TCommonBlock) stmt).getBlockBody().getParentObjectName() instanceof TStoredProcedureSqlStatement) {
506                    return (TStoredProcedureSqlStatement) ((TCommonBlock) stmt).getBlockBody().getParentObjectName();
507                }
508                else {
509                        return getParentProcedure(stmt.getParentStmt());
510                }
511            }
512            
513            return null;
514        }
515        
516            stmt = stmt.getParentStmt();
517        if (stmt == null) {
518            if(ModelBindingManager.getGlobalProcedure()!=null) {
519                return ModelBindingManager.getGlobalProcedure();
520            }
521            else {
522                return null;
523            }
524        }
525        
526        if (stmt instanceof TCommonBlock) {
527            if (((TCommonBlock) stmt).getBlockBody().getParentObjectName() instanceof TStoredProcedureSqlStatement) {
528                stmt = (TStoredProcedureSqlStatement) ((TCommonBlock) stmt).getBlockBody().getParentObjectName();
529            }
530        }
531
532        return getParentProcedure(stmt);
533    }
534
535    public Pair3<Long, Long, String> getStartPosition() {
536        return startPosition;
537    }
538
539    public Pair3<Long, Long, String> getEndPosition() {
540        return endPosition;
541    }
542
543    public TCustomSqlStatement getGspObject() {
544        return gspObject;
545    }
546
547    public long getId() {
548        return id;
549    }
550
551    public String getSchema() {
552        return schema;
553    }
554
555    public String getDatabase() {
556        return database;
557    }
558
559    public String getName() {
560        return name;
561    }
562
563    public String getType() {
564        return type;
565    }
566
567    public void setName(String name) {
568        this.name = name;
569    }
570
571    public void setType(String type) {
572        this.type = type;
573    }
574
575    public String getProcedureName() {
576        return procedureName;
577    }
578
579    public void setProcedureName(String procedureName) {
580        this.procedureName = procedureName;
581    }
582
583    public String getQueryHashId() {
584        return queryHashId;
585    }
586
587    public void setQueryHashId(String queryHashId) {
588        this.queryHashId = queryHashId;
589    }
590
591    public void appendColumn(Object tableColumn) {
592        if (tableColumn != null && !targetColumns.contains(tableColumn)) {
593            targetColumns.add(tableColumn);
594        }
595    }
596
597    public List<Object> getTargetColumns() {
598        return targetColumns;
599    }
600
601    public String getCustomType() {
602        return customType;
603    }
604
605    public Long getProcedureId() {
606        return procedureId;
607    }
608
609    public String getServer() {
610        return server;
611    }
612
613    public List<Transform> getTransforms() {
614        return transforms;
615    }
616}