| 32 | { |
| 33 | |
| 34 | HostWithPorts getTargetServer(ContextPtr context, ASTPtr & ast) |
| 35 | { |
| 36 | auto get_target_server_for_table = [&] (const String & database_name, const String & table_name) -> HostWithPorts |
| 37 | { |
| 38 | if (database_name == "system") |
| 39 | return {}; |
| 40 | |
| 41 | DatabaseAndTable db_and_tb = DatabaseCatalog::instance().tryGetDatabaseAndTable(StorageID(database_name, table_name), context); |
| 42 | DatabasePtr db_ptr = std::move(db_and_tb.first); |
| 43 | StoragePtr storage_ptr = std::move(db_and_tb.second); |
| 44 | if (!db_ptr || !storage_ptr || db_ptr->getEngineName() != "Cnch") |
| 45 | return {}; |
| 46 | |
| 47 | auto topology_master = context->getCnchTopologyMaster(); |
| 48 | |
| 49 | return topology_master->getTargetServer( |
| 50 | UUIDHelpers::UUIDToString(storage_ptr->getStorageUUID()), storage_ptr->getServerVwName(), context->getTimestamp(), true); |
| 51 | }; |
| 52 | |
| 53 | if (const auto * alter = ast->as<ASTAlterQuery>()) |
| 54 | { |
| 55 | if (alter->alter_object == ASTAlterQuery::AlterObjectType::DATABASE) |
| 56 | return {}; |
| 57 | return get_target_server_for_table(alter->database.empty() ? context->getCurrentDatabase() : alter->database, alter->table); |
| 58 | } |
| 59 | else if (const auto * alter_mysql = ast->as<ASTAlterAnalyticalMySQLQuery>()) |
| 60 | { |
| 61 | if(alter_mysql->alter_object == ASTAlterQuery::AlterObjectType::DATABASE) |
| 62 | return {}; |
| 63 | return get_target_server_for_table(alter_mysql->database.empty() ? context->getCurrentDatabase() : alter_mysql->database, alter_mysql->table); |
| 64 | } |
| 65 | else if (const auto * select = ast->as<ASTSelectWithUnionQuery>()) |
| 66 | { |
| 67 | if (!context->getSettingsRef().enable_select_query_forwarding && !context->getSettingsRef().enable_multiple_table_select_query_forwarding) |
| 68 | return {}; |
| 69 | |
| 70 | ASTs tables; |
| 71 | bool has_table_func = false; |
| 72 | ASTSelectQuery::collectAllTables(ast.get(), tables, has_table_func); |
| 73 | |
| 74 | if (tables.empty() || has_table_func) |
| 75 | return {}; |
| 76 | |
| 77 | String current_db = context->getCurrentDatabase(); |
| 78 | if (tables.size() == 1) |
| 79 | { |
| 80 | DatabaseAndTableWithAlias db_and_table(tables[0], current_db); |
| 81 | LOG_DEBUG( |
| 82 | &Poco::Logger::get("executeQuery"), |
| 83 | "Get main table `{}.{}` for current select query.", |
| 84 | db_and_table.database, |
| 85 | db_and_table.table); |
| 86 | return get_target_server_for_table(db_and_table.database, db_and_table.table); |
| 87 | } |
| 88 | else |
| 89 | { |
| 90 | if (!context->getSettingsRef().enable_multiple_table_select_query_forwarding) |
| 91 | return {}; |
no test coverage detected