YoushouldhavereceivedacopyoftheGNUGeneralPublicLicense alongwiththisprogram;ifnot,writetotheFreeSoftware
Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1335 USA */
ulong wsrep_debug; // Debug level logging
my_bool wsrep_convert_LOCK_to_trx; // Convert locking sessions to trx
my_bool wsrep_auto_increment_control; // Control auto increment variables
my_bool wsrep_drupal_282555_workaround; // Retry autoinc insert after dupkey
my_bool wsrep_certify_nonPK; // Certify, even when no primary key
ulong wsrep_certification_rules = WSREP_CERTIFICATION_RULES_STRICT;
READ_ONLY_SYSVAR my_bool wsrep_recovery; // Recovery
my_bool wsrep_log_conflicts;
my_bool wsrep_load_data_splitting= 0; // Commit load data every 10K intervals
my_bool wsrep_slave_UK_checks; // Slave thread does UK checks
my_bool wsrep_slave_FK_checks; // Slave thread does FK checks
my_bool wsrep_restart_slave; // Should mysql slave thread be // restarted, when node joins back?
my_bool wsrep_desync; // De(re)synchronize the node from the // cluster
ulonglong wsrep_mode; bool wsrep_service_started; // If Galera was initialized long wsrep_slave_threads; // No. of slave appliers threads
ulong wsrep_retry_autocommit; // Retry aborted autocommit trx
ulong wsrep_max_ws_size; // Max allowed ws (RBR buffer) size
ulong wsrep_max_ws_rows; // Max number of rows in ws
ulong wsrep_forced_binlog_format= BINLOG_FORMAT_UNSPEC;
ulong wsrep_mysql_replication_bundle;
bool wsrep_gtid_mode; // Enable WSREP native GTID support
Wsrep_gtid_server wsrep_gtid_server;
uint wsrep_gtid_domain_id=0; // Domain id on above structure
/* Other configuration variables and their default values. */
my_bool wsrep_incremental_data_collection= 0; // Incremental data collection bool wsrep_new_cluster= false; // Bootstrap the cluster? int wsrep_slave_count_change= 0; // No. of appliers to stop/start int wsrep_to_isolation= 0; // No. of active TO isolation threads // // NOTE : MySQL Cluster has max_protocol version 7 // thus protocol versions 5-7 are reserved for compatibility. // long wsrep_max_protocol_version= 4; // Maximum protocol version to use longint wsrep_protocol_version= wsrep_max_protocol_version;
ulong wsrep_trx_fragment_unit= WSREP_FRAG_BYTES; // unit for fragment size
READ_ONLY_SYSVAR ulong wsrep_SR_store_type= WSREP_SR_STORE_TABLE;
uint wsrep_ignore_apply_errors= 0;
uint wsrep_applier_retry_count= 0;
std::atomic <bool> wsrep_thread_create_failed;
/* *Endconfigurationoptions
*/
/* *Cachedvariables
*/
// Whether the Galera write-set replication provider is set // wsrep_provider && strcmp(wsrep_provider, WSREP_NONE) bool WSREP_PROVIDER_EXISTS_;
// Whether the Galera write-set replication is enabled // global_system_variables.wsrep_on && WSREP_PROVIDER_EXISTS_ bool WSREP_ON_;
switch (level) { case wsrep::log::info:
WSREP_INFO("%s", msg); break; case wsrep::log::warning:
WSREP_WARN("%s", msg); break; case wsrep::log::error:
WSREP_ERROR("%s", msg); break; case wsrep::log::debug:
WSREP_DEBUG("%s", msg); break; case wsrep::log::unknown:
WSREP_UNKNOWN("%s", msg); break;
}
}
void wsrep_init_gtid()
{
wsrep_server_gtid_t stored_gtid= wsrep_get_SE_checkpoint<wsrep_server_gtid_t>(); // Domain id may have changed, use the one // received during state transfer.
stored_gtid.domain_id= wsrep_gtid_server.domain_id; if (stored_gtid.server_id == 0)
{
rpl_gtid wsrep_last_gtid; if (mysql_bin_log.is_open() &&
mysql_bin_log.lookup_domain_in_binlog_state(stored_gtid.domain_id,
&wsrep_last_gtid))
{
stored_gtid.server_id= wsrep_last_gtid.server_id;
stored_gtid.seqno= wsrep_last_gtid.seq_no;
} else
{
stored_gtid.server_id= global_system_variables.server_id;
stored_gtid.seqno= 0;
}
}
wsrep_gtid_server.gtid(stored_gtid);
}
WSREP_INFO("wsrep_init_schema_and_SR %p", wsrep_schema); if (!wsrep_schema)
{
wsrep_schema= new Wsrep_schema(); if (wsrep_schema->init())
{
WSREP_ERROR("Failed to init wsrep schema");
unireg_abort(1);
} // If we are bootstraping new cluster we should // clear allowlist table and populate it from variable if (wsrep_new_cluster)
{
wsrep_schema->clear_allowlist();
std::vector<std::string> ip_allowlist; if (wsrep_split_allowlist(ip_allowlist))
{
wsrep_schema->store_allowlist(ip_allowlist);
}
}
}
}
void wsrep_recover_sr_from_storage(THD *orig_thd)
{ switch (wsrep_SR_store_type)
{ case WSREP_SR_STORE_TABLE: if (!wsrep_schema)
{
WSREP_ERROR("Wsrep schema not initialized when trying to recover " "streaming transactions: wsrep_on %d", WSREP_ON);
trans_commit(orig_thd);
} if (wsrep_schema->recover_sr_transactions(orig_thd))
{
WSREP_ERROR("Failed to recover SR transactions from schema: wsrep_on : %d", WSREP_ON);
trans_commit(orig_thd);
} break; default: /* */
WSREP_ERROR("Unsupported wsrep SR store type: %lu wsrep_on: %d",
wsrep_SR_store_type, WSREP_ON);
trans_commit(orig_thd); break;
}
}
/** Export the WSREP provider's capabilities as a human readable string. *Theresultissavedinadynamicallyallocatedstringoftheform: *:cap1:cap2:cap3:
*/ staticvoid wsrep_capabilities_export(wsrep_cap_t const cap, char** str)
{ staticconstchar* names[] =
{ /* Keep in sync with wsrep/wsrep_api.h WSREP_CAP_* macros. */ "MULTI_MASTER", "CERTIFICATION", "PARALLEL_APPLYING", "TRX_REPLAY", "ISOLATION", "PAUSE", "CAUSAL_READS", "CAUSAL_TRX", "INCREMENTAL_WRITESET", "SESSION_LOCKS", "DISTRIBUTED_LOCKS", "CONSISTENCY_CHECK", "UNORDERED", "ANNOTATION", "PREORDERED", "STREAMING", "SNAPSHOT", "NBO",
};
std::string s; for (size_t i= 0; i < sizeof(names) / sizeof(names[0]); ++i)
{ if (cap & (1ULL << i))
{ if (s.empty())
{
s= ":";
}
s += names[i];
s += ":";
}
}
/* A read from the string pointed to by *str may be started at any time,
* so it must never point to free(3)d memory or non '\0' terminated string. */
char* const previous= *str;
*str= strdup(s.c_str());
if (previous != NULL)
{
free(previous);
}
}
/* Verifies that SE position is consistent with the group position
* and initializes other variables */ void wsrep_verify_SE_checkpoint(const wsrep_uuid_t& uuid,
wsrep_seqno_t const seqno)
{
}
Wsrep_server_state::init_once(server_name,
incoming_address,
node_address,
working_dir,
initial_position,
wsrep_max_protocol_version);
Wsrep_server_state::instance().debug_log_level(wsrep_debug);
} catch (const wsrep::runtime_error& e)
{
WSREP_ERROR("Failed to init wsrep server %s", e.what()); return1;
} catch (const std::exception& e)
{
WSREP_ERROR("Failed to init wsrep server %s", e.what());
} return0;
}
void wsrep_init_globals()
{
wsrep_init_sidno(Wsrep_server_state::instance().connected_gtid().id()); /* Recover last written wsrep gtid */
wsrep_init_gtid(); if (wsrep_new_cluster)
{ /* Start with provided domain_id & server_id found in configuration */
wsrep_server_gtid_t new_gtid;
new_gtid.domain_id= wsrep_gtid_domain_id;
new_gtid.server_id= global_system_variables.server_id; /* Use seqno which was recovered in wsrep_init_gtid() */
new_gtid.seqno= wsrep_gtid_server.seqno(); /* Try to search for domain_id and server_id combination in binlog if found continue from last seqno */
wsrep_get_binlog_gtid_seqno(new_gtid);
wsrep_gtid_server.gtid(new_gtid);
} else
{ if (wsrep_gtid_mode && wsrep_gtid_server.server_id != global_system_variables.server_id)
WSREP_INFO("Ignoring server id %ld for non bootstrap node, using %ld.",
global_system_variables.server_id, wsrep_gtid_server.server_id);
}
wsrep_init_schema();
{ /* apparently this thread has already called my_thread_init(),
* so we skip it, hence 'false' for initialization. */
wsp::thd thd(false, true);
thd.ptr->variables.tx_read_only= thd.ptr->tx_read_only= false;
wsrep_sst_cleanup_user(thd.ptr);
} if (WSREP_ON)
{
Wsrep_server_state::instance().initialized();
}
}
if (!*wsrep_provider ||
!strcasecmp(wsrep_provider, WSREP_NONE))
{ // enable normal operation in case no provider is specified
global_system_variables.wsrep_on= 0; int err= Wsrep_server_state::init_provider(
wsrep_provider, wsrep_provider_options ? wsrep_provider_options : ""); if (err)
{
DBUG_PRINT("wsrep",("wsrep::init() failed: %d", err));
WSREP_ERROR("wsrep::init() failed: %d, must shutdown", err);
} else
wsrep_init_provider_status_variables(); return err;
}
if (wsrep_gtid_mode && opt_bin_log && !opt_log_slave_updates)
{
WSREP_ERROR("Option --log-slave-updates is required if " "binlog is enabled, GTID mode is on and wsrep provider " "is specified"); return1;
}
if (!wsrep_data_home_dir || strlen(wsrep_data_home_dir) == 0)
wsrep_data_home_dir= mysql_real_data_home;
/* In case of errors, SST tmp dir is not set */
wsrep_sst_tmp_dir_check();
Wsrep_server_state::init_provider_services(); if (Wsrep_server_state::instance().load_provider(
wsrep_provider,
wsrep_provider_options,
Wsrep_server_state::instance().provider_services()))
{
WSREP_ERROR("Failed to load provider");
Wsrep_server_state::deinit_provider_services(); return1;
}
if (!wsrep_provider_is_SR_capable() &&
global_system_variables.wsrep_trx_fragment_size > 0)
{
WSREP_ERROR("The WSREP provider (%s) does not support streaming " "replication but wsrep_trx_fragment_size is set to a " "value other than 0 (%llu). Cannot continue. Either set " "wsrep_trx_fragment_size to 0 or use wsrep_provider that " "supports streaming replication.",
wsrep_provider, global_system_variables.wsrep_trx_fragment_size);
Wsrep_server_state::instance().deinit_provider();
Wsrep_server_state::deinit_provider_services(); return1;
}
/* Now WSREP is fully initialized */
global_system_variables.wsrep_on= 1;
WSREP_ON_= wsrep_provider && *wsrep_provider && strcasecmp(wsrep_provider, WSREP_NONE);
wsrep_service_started= 1;
staticvoid wsrep_stop_replication_common(THD *thd)
{ if (Wsrep_server_state::instance().state() !=
Wsrep_server_state::s_disconnected)
{
WSREP_DEBUG("Disconnect provider");
Wsrep_server_state::instance().disconnect(); if (Wsrep_server_state::instance().wait_until_state(
Wsrep_server_state::s_disconnected))
{
WSREP_WARN("Wsrep interrupted while waiting for disconnected state");
}
}
/* my connection, should not terminate with
wsrep_close_client_connections(), make transaction to rollback */ if (thd && !thd->wsrep_applier)
trans_rollback(thd);
wsrep_close_client_connections(TRUE, thd);
/* wait until appliers have stopped */
wsrep_wait_appliers_close(thd);
void wsrep_shutdown()
{ /* Signal ready state waiters that we're shutting down. */
mysql_mutex_lock(&LOCK_wsrep_ready);
DBUG_ASSERT(wsrep_state != IN_SHUTDOWN);
wsrep_state = IN_SHUTDOWN;
mysql_cond_signal(&COND_wsrep_ready);
mysql_mutex_unlock(&LOCK_wsrep_ready);
/* Stop wsrep threads in case they are running. */ if (wsrep_running_threads > 0)
{
WSREP_INFO("Shutdown replication");
wsrep_stop_replication_common(nullptr); /* Undocking the thread specific data. */
set_current_thd(nullptr);
}
}
bool wsrep_start_replication(constchar *wsrep_cluster_address)
{ int rcode;
WSREP_DEBUG("wsrep_start_replication");
/* ifprovideristrivial,don'teventrytoconnect, butresumelocalnodeoperation
*/ if (!WSREP_PROVIDER_EXISTS)
{ // enable normal operation in case no provider is specified returntrue;
}
DBUG_ASSERT(wsrep_cluster_address[0]);
// --wsrep-new-cluster flag is not used, checking wsrep_cluster_address // it should match gcomm:// only to be considered as bootstrap node. // This logic is used in galera. if (!wsrep_new_cluster &&
(strlen(wsrep_cluster_address) == 8) &&
!strncmp(wsrep_cluster_address, "gcomm://", 8))
{
wsrep_new_cluster= true;
}
//seconds after which the limit warnings suppression will be activated #define WSREP_WARNING_ACTIVATION_TIMEOUT 5*60 //number of limit warnings after which the suppression will be activated #define WSREP_WARNING_ACTIVATION_THRESHOLD 10
if (!wsrep_warning_active[warning_type])
{ /* ACTIVATION: WegotWSREP_WARNING_ACTIVATION_THRESHOLDwarningsin lessthanWSREP_WARNING_ACTIVATION_TIMEOUTweactivatethe suppression.
*/ if (diff_time <= WSREP_WARNING_ACTIVATION_TIMEOUT)
{
wsrep_warning_active[warning_type]= true;
WSREP_INFO("Suppressing warnings of type '%s' for up to %d seconds because of flooding",
wsrep_warning_name(warning_type),
WSREP_WARNING_ACTIVATION_TIMEOUT);
} else
{ /* Thereisnofloodingtillnow,thereforewerestartthemonitoring
*/
wsrep_reset_warnings(now);
}
} else
{ /* This type of warnings was suppressed */ if (diff_time > WSREP_WARNING_ACTIVATION_TIMEOUT)
{
ulonglong save_count= wsrep_total_warnings_count; /* Print a suppression note and remove the suppression */
wsrep_reset_warnings(now);
WSREP_INFO("Suppressed %lu unsafe warnings during " "the last %d seconds",
save_count, (int) diff_time);
}
}
}
return wsrep_warning_active[warning_type];
}
/** Auxiliaryfunctiontopushwarningtoclientandtotheerrorlog
*/ staticvoid wsrep_push_warning(THD *thd, enum wsrep_warning_type type, const handlerton *hton, const TABLE_LIST *tables)
{ switch(type)
{ case WSREP_REQUIRE_PRIMARY_KEY:
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
ER_OPTION_PREVENTS_STATEMENT, "WSREP: wsrep_mode = REQUIRED_PRIMARY_KEY enabled. " "Table '%s'.'%s' should have PRIMARY KEY defined.",
tables->db.str, tables->table_name.str); if (global_system_variables.log_warnings > 1 &&
!wsrep_protect_against_warning_flood(type))
WSREP_WARN("wsrep_mode = REQUIRED_PRIMARY_KEY enabled. " "Table '%s'.'%s' should have PRIMARY KEY defined",
tables->db.str, tables->table_name.str); break; case WSREP_REQUIRE_INNODB:
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
ER_OPTION_PREVENTS_STATEMENT, "WSREP: wsrep_mode = STRICT_REPLICATION enabled. " "Storage engine %s for table '%s'.'%s' is " "not supported in Galera",
ha_resolve_storage_engine_name(hton),
tables->db.str, tables->table_name.str); if (global_system_variables.log_warnings > 1 &&
!wsrep_protect_against_warning_flood(type))
WSREP_WARN("wsrep_mode = STRICT_REPLICATION enabled. " "Storage engine %s for table '%s'.'%s' is " "not supported in Galera",
ha_resolve_storage_engine_name(hton),
tables->db.str, tables->table_name.str); break; case WSREP_EXPERIMENTAL:
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
ER_ERROR_DURING_COMMIT, "WSREP: Replication of non-transactional engines is experimental. " "Storage engine %s for table '%s'.'%s' is " "not supported in Galera",
ha_resolve_storage_engine_name(hton),
tables->db.str, tables->table_name.str); if (global_system_variables.log_warnings > 1 &&
!wsrep_protect_against_warning_flood(type))
WSREP_WARN("Replication of non-transactional engines is experimental. " "Storage engine %s for table '%s'.'%s' is " "not supported in Galera",
ha_resolve_storage_engine_name(hton),
tables->db.str, tables->table_name.str); break; default: assert(0); break;
}
}
TABLE *tbl= tables->table; /* If this is partitioned table we need to find out implementingstorageenginehandlerton.
*/ const handlerton *ht= tbl->file->partition_ht(); if (!ht) ht= hton;
DBUG_ASSERT(tbl); if (replicate)
{ /* ItisnotrecommendedtoreplicateMyISAMasitlacksrollback featurebutifuserdemandsthenactionsarereplicatedusingTOI. Followingcodewillkick-starttheTOIbutthishastobedone onlyonceperstatement.
Note:kick-startwilltakecareofcreatingisolationkeyfor alltablesinvolvedinthelist(providedallofthemareMYISAM orAriatables).
*/ if (tbl->s->table_category != TABLE_CATEGORY_STATISTICS)
{ if (tbl->s->primary_key == MAX_KEY &&
wsrep_check_mode(WSREP_MODE_REQUIRED_PRIMARY_KEY))
{ /* Other replicated table doesn't have explicit primary-key defined. */
wsrep_push_warning(thd, WSREP_REQUIRE_PRIMARY_KEY, hton, tables);
}
if (wsrep_check_mode(WSREP_MODE_STRICT_REPLICATION))
{ /* Table is not an InnoDB table and strict replication is requested*/
wsrep_push_warning(thd, WSREP_REQUIRE_INNODB, hton, tables);
}
// Check are we inside a transaction constbool changes= wsrep_has_changes(thd); constbool active= wsrep_is_active(thd);
// We should not start TOI if transaction has made already // changes and is active if (changes && active)
{
my_message(ER_ERROR_DURING_COMMIT, "Transactional commit not supported " "by involved engine(s)", MYF(0));
wsrep_push_warning(thd, WSREP_EXPERIMENTAL, hton, tables); returnfalse;
}
// Roll back current stmt if exists
wsrep_before_rollback(thd, true);
wsrep_after_rollback(thd, true);
wsrep_after_statement(thd);
if (!is_system_db &&
!is_temporary_table(tables))
{ if (db_type != DB_TYPE_INNODB &&
wsrep_check_mode(WSREP_MODE_STRICT_REPLICATION))
{ /* Table is not an InnoDB table and strict replication is requested*/
wsrep_push_warning(thd, WSREP_REQUIRE_INNODB, hton, tables);
}
if (db_type != DB_TYPE_INNODB &&
thd->variables.sql_log_bin == 1 &&
wsrep_check_mode(WSREP_MODE_DISALLOW_LOCAL_GTID))
{ /* Table is not an InnoDB table and local GTIDs are disallowed */
my_error(ER_GALERA_REPLICATION_NOT_SUPPORTED, MYF(0));
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
ER_OPTION_PREVENTS_STATEMENT, "You can't execute statements that would generate local " "GTIDs when wsrep_mode = DISALLOW_LOCAL_GTID is set. " "Try disabling binary logging with SET sql_log_bin=0 " "to execute this statement."); goto wsrep_error_label;
}
}
}
returntrue;
wsrep_error_label: returnfalse;
}
bool wsrep_check_mode_before_cmd_execute (THD *thd)
{ bool ret= true; if (wsrep_check_mode(WSREP_MODE_BINLOG_ROW_FORMAT_ONLY) &&
!thd->is_current_stmt_binlog_format_row() && is_update_query(thd->lex->sql_command))
{
my_error(ER_GALERA_REPLICATION_NOT_SUPPORTED, MYF(0));
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
ER_OPTION_PREVENTS_STATEMENT, "WSREP: wsrep_mode = BINLOG_ROW_FORMAT_ONLY enabled. Only ROW binlog format is supported.");
ret= false;
} if (wsrep_check_mode(WSREP_MODE_REQUIRED_PRIMARY_KEY) &&
thd->lex->sql_command == SQLCOM_CREATE_TABLE)
{
Key *key;
List_iterator<Key> key_iterator(thd->lex->alter_info.key_list); bool primary_key_found= false; while ((key= key_iterator++))
{ if (key->type == Key::PRIMARY)
{
primary_key_found= true; break;
}
} if (!primary_key_found)
{
my_error(ER_GALERA_REPLICATION_NOT_SUPPORTED, MYF(0));
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
ER_OPTION_PREVENTS_STATEMENT, "WSREP: wsrep_mode = REQUIRED_PRIMARY_KEY enabled. Table should have PRIMARY KEY defined.");
ret= false;
}
} return ret;
}
staticenum enum_wsrep_sync_wait
wsrep_sync_wait_mask_for_command(enum enum_sql_command command)
{ switch (command)
{ case SQLCOM_SELECT: case SQLCOM_CHECKSUM: return WSREP_SYNC_WAIT_BEFORE_READ; case SQLCOM_DELETE: case SQLCOM_DELETE_MULTI: case SQLCOM_UPDATE: case SQLCOM_UPDATE_MULTI: return WSREP_SYNC_WAIT_BEFORE_UPDATE_DELETE; case SQLCOM_REPLACE: case SQLCOM_INSERT: case SQLCOM_REPLACE_SELECT: case SQLCOM_INSERT_SELECT: return WSREP_SYNC_WAIT_BEFORE_INSERT_REPLACE; default: if (wsrep_is_diagnostic_query(command))
{ return WSREP_SYNC_WAIT_NONE;
} if (wsrep_is_show_query(command))
{ switch (command)
{ case SQLCOM_SHOW_PROFILE: case SQLCOM_SHOW_PROFILES: case SQLCOM_SHOW_SLAVE_HOSTS: case SQLCOM_SHOW_RELAYLOG_EVENTS: case SQLCOM_SHOW_SLAVE_STAT: case SQLCOM_SHOW_BINLOG_STAT: case SQLCOM_SHOW_ENGINE_STATUS: case SQLCOM_SHOW_ENGINE_MUTEX: case SQLCOM_SHOW_ENGINE_LOGS: case SQLCOM_SHOW_PROCESSLIST: case SQLCOM_SHOW_PRIVILEGES: return WSREP_SYNC_WAIT_NONE; default: return WSREP_SYNC_WAIT_BEFORE_SHOW;
}
}
} return WSREP_SYNC_WAIT_NONE;
}
bool wsrep_sync_wait(THD* thd, enum enum_sql_command command)
{ bool res = false; if (WSREP_CLIENT(thd) && thd->variables.wsrep_sync_wait)
res = wsrep_sync_wait(thd, wsrep_sync_wait_mask_for_command(command)); return res;
}
/* if there is prepare query, add event for it */ if (!ret && thd->wsrep_TOI_pre_query)
{
Query_log_event ev(thd, thd->wsrep_TOI_pre_query,
thd->wsrep_TOI_pre_query_len, FALSE, FALSE, FALSE, 0); if (writer.write(&ev)) ret= 1;
}
/* continue to append the actual query */
Query_log_event ev(thd, query, query_len, FALSE, FALSE, FALSE, 0); /* WSREP GTID mode, we need to change server_id */ if (wsrep_gtid_mode && !thd->variables.gtid_seq_no)
ev.server_id= wsrep_gtid_server.server_id; if (!ret && writer.write(&ev)) ret= 1; if (!ret && wsrep_write_cache_buf(&tmp_io_cache, buf, buf_len)) ret= 1;
close_cached_file(&tmp_io_cache); return ret;
}
staticint
wsrep_alter_query_string(THD *thd, String *buf)
{ /* Append the "ALTER" part of the query */ if (buf->append(STRING_WITH_LEN("ALTER "))) return1; /* Append definer */
append_definer(thd, buf, &(thd->lex->definer->user), &(thd->lex->definer->host)); /* Append the left part of thd->query after event name part */ if (buf->append(thd->lex->stmt_definition_begin,
thd->lex->stmt_definition_end -
thd->lex->stmt_definition_begin)) return1;
if (definer)
{
views->definer.user= definer->user;
views->definer.host= definer->host;
} else {
WSREP_ERROR("Failed to get DEFINER for VIEW."); return1;
}
view_store_options(thd, views, &buff);
buff.append(STRING_WITH_LEN("VIEW ")); /* Test if user supplied a db (ie: we did not use thd->db) */ if (views->db.str && views->db.str[0] &&
(thd->db.str == NULL || cmp(&views->db, &thd->db)))
{
append_identifier(thd, &buff, &views->db);
buff.append('.');
}
append_identifier(thd, &buff, &views->table_name); if (lex->view_list.elements)
{
List_iterator_fast<LEX_CSTRING> names(lex->view_list);
LEX_CSTRING *name; int i;
/*! Should DDL be replicated by Galera * *@paramthdthreadhandle *@paramhtonrealstorageenginehandlerton *
* @retval true if we should replicate DDL, false if not */
bool wsrep_should_replicate_ddl(THD* thd, const handlerton *hton)
{ if (!wsrep_check_mode(WSREP_MODE_STRICT_REPLICATION)) returntrue;
DBUG_ASSERT(hton != nullptr);
switch (hton->db_type)
{ case DB_TYPE_UNKNOWN: /* Special pseudo-handlertons (such as 10.6+ JSON tables). */ returntrue; break; case DB_TYPE_INNODB: returntrue; break; case DB_TYPE_MYISAM: if (wsrep_check_mode(WSREP_MODE_REPLICATE_MYISAM)) returntrue; break; case DB_TYPE_ARIA: if (wsrep_check_mode(WSREP_MODE_REPLICATE_ARIA)) returntrue; break; case DB_TYPE_PARTITION_DB: /* In most cases this means we could not find out
table->file->partition_ht() */ returntrue; break; default: break;
}
WSREP_DEBUG("wsrep OSU failed for %s", wsrep_thd_query(thd));
bool wsrep_should_replicate_ddl_iterate(THD* thd, const TABLE_LIST* table_list)
{ for (const TABLE_LIST* it= table_list; it; it= it->next_global)
{ const TABLE* table= it->table; if (table && !it->table_function)
{ /* If this is partitioned table we need to find out implementingstorageenginehandlerton.
*/ const handlerton *ht= table->file->partition_ht(); if (!ht) ht= table->s->db_type(); if (!wsrep_should_replicate_ddl(thd, ht)) returnfalse;
}
} returntrue;
}
switch (lex->sql_command)
{ case SQLCOM_CREATE_TABLE: if (thd->lex->create_info.options & HA_LEX_CREATE_TMP_TABLE)
{ returnfalse;
} if (!wsrep_should_replicate_ddl(thd, create_info->db_type))
{ returnfalse;
} /* IfmariadbmasterhasreplicatedaCTAS,weshouldnotreplicatethecreatetable partseparatelyasTOI,buttoreplicatebothcreatetableandfollowinginserts asonewriteset. However,ifCTAScreatesemptytable,weshouldreplicatethecreatetablealone asTOI.Wehavetodorelaylogeventlookuptoseeifroweventsfollowthe createtableevent.
*/ if (thd->slave_thread &&
!(thd->rgi_slave->gtid_ev_flags2 & Gtid_log_event::FL_STANDALONE))
{ /* this is CTAS, either empty or populated table */
ulonglong event_size = 0; enum Log_event_type ev_type= wsrep_peak_event(thd->rgi_slave, &event_size); switch (ev_type)
{ case QUERY_EVENT: case XID_EVENT: /* CTAS with empty table, we replicate create table as TOI */ break;
case TABLE_MAP_EVENT:
WSREP_DEBUG("replicating CTAS of empty table as TOI"); // fall through case WRITE_ROWS_EVENT: /* CTAS with populated table, we replicate later at commit time */
WSREP_DEBUG("skipping create table of CTAS replication"); returnfalse;
default:
WSREP_WARN("unexpected async replication event: %d", ev_type);
} returntrue;
} /* no next async replication event */ returntrue; break; case SQLCOM_CREATE_VIEW:
DBUG_ASSERT(!table_list);
DBUG_ASSERT(first_table); /* First table is view name */ /* Ifanyoftheremainingtablesrefertotemporarytableerror isreturnedtoclient,soTOIcanbeskipped
*/ for (const TABLE_LIST* it= first_table->next_global; it; it= it->next_global)
{ if (thd->find_temporary_table(it))
{ returnfalse;
}
} returntrue; break; case SQLCOM_CREATE_TRIGGER:
DBUG_ASSERT(first_table);
if (thd->find_temporary_table(first_table))
{ returnfalse;
} returntrue; break; case SQLCOM_DROP_TRIGGER:
DBUG_ASSERT(table_list); if (thd->find_temporary_table(table_list))
{ returnfalse;
} returntrue; break; case SQLCOM_ALTER_TABLE: if (create_info)
{ const handlerton *hton= create_info->db_type; if (!hton)
hton= ha_default_handlerton(thd); if (!wsrep_should_replicate_ddl(thd, hton)) returnfalse;
} /* fallthrough */ default: if (table && !thd->find_temporary_table(Lex_ident_db(Lex_cstring_strlen(db)),
Lex_cstring_strlen(table)))
{ returntrue;
}
if (table_list)
{ for (const TABLE_LIST* table= first_table; table; table= table->next_global)
{ if (!thd->find_temporary_table(table->db, table->table_name))
{ returntrue;
}
}
}
return !(table || table_list); break; case SQLCOM_CREATE_SEQUENCE: /* No TOI for temporary sequences as they are
not replicated */ if (thd->lex->tmp_table())
{ returnfalse;
} returntrue;
if (buf_err) {
WSREP_ERROR("Failed to create TOI event buf: %d", buf_err);
my_message(ER_UNKNOWN_ERROR, "WSREP replication failed to prepare TOI event buffer. " "Check your query.",
MYF(0)); return -1;
}
if (thd->has_read_only_protection())
{ /* non replicated DDL, affecting temporary tables only */
WSREP_DEBUG("TO isolation skipped, sql: %s." "Only temporary tables affected.",
wsrep_thd_query(thd)); if (buf) my_free(buf); return -1;
}
thd_proc_info(thd, "acquiring total order isolation");
WSREP_DEBUG("wsrep_TOI_begin for %s", wsrep_thd_query(thd));
THD_STAGE_INFO(thd, stage_waiting_isolation);
DEBUG_SYNC(thd, "wsrep_before_toi_begin");
wsrep::client_state& cs(thd->wsrep_cs());
int ret= cs.enter_toi_local(key_array,
wsrep::const_buffer(buff.ptr, buff.len));
if (ret)
{
DBUG_ASSERT(cs.current_error());
WSREP_WARN("TO isolation error %s for : %s.%s, sql: %s. ",
wsrep::to_c_string(cs.current_error()),
(db ? db : "(null)"),
(table ? table : " "),
wsrep_thd_query(thd));
/* jump to error handler in mysql_execute_command() */ switch (cs.current_error())
{ case wsrep::e_size_exceeded_error:
my_error(ER_UNKNOWN_ERROR, MYF(0), "Maximum writeset size exceeded"); break; case wsrep::e_deadlock_error:
my_error(ER_LOCK_DEADLOCK, MYF(0)); break; case wsrep::e_timeout_error:
my_error(ER_LOCK_WAIT_TIMEOUT, MYF(0)); break; default: if (!thd->is_error())
{
my_error(ER_LOCK_DEADLOCK, MYF(0));
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
ER_LOCK_DEADLOCK, "WSREP replication failed with error %s. " "Check your wsrep connection state and retry the query.",
wsrep::to_c_string(cs.current_error()));
}
}
rc= -1;
} else { if (!thd->variables.gtid_seq_no)
{
uint64 seqno= 0; if (thd->variables.wsrep_gtid_seq_no &&
thd->variables.wsrep_gtid_seq_no > wsrep_gtid_server.seqno())
{
seqno= thd->variables.wsrep_gtid_seq_no;
wsrep_gtid_server.seqno(thd->variables.wsrep_gtid_seq_no);
} else
{
seqno= wsrep_gtid_server.seqno_inc();
}
thd->variables.wsrep_gtid_seq_no= 0;
thd->wsrep_current_gtid_seqno= seqno; if (mysql_bin_log.is_open() && wsrep_gtid_mode)
{
thd->variables.gtid_seq_no= seqno;
thd->variables.gtid_domain_id= wsrep_gtid_server.domain_id;
thd->variables.server_id= wsrep_gtid_server.server_id;
}
}
++wsrep_to_isolation;
rc= 0;
}
if (thd->is_error() && !wsrep_must_ignore_error(thd))
{ /* use only error code, for the message can be inconsistent *betweenthenodesduetodifferinglc_messagesettings
* in client session and server applier thread */
wsrep_store_error(thd, err, false);
}
/* For CREATE TEMPORARY SEQUENCE we do not start RSU because objectislocalonlyandactuallyCREATETABLE+INSERT
*/ if (thd->lex->sql_command == SQLCOM_CREATE_SEQUENCE &&
thd->lex->tmp_table()) return1;
if (thd->variables.wsrep_OSU_method == WSREP_OSU_RSU &&
thd->variables.sql_log_bin == 1 &&
wsrep_check_mode(WSREP_MODE_DISALLOW_LOCAL_GTID))
{ /* wsrep_mode = WSREP_MODE_DISALLOW_LOCAL_GTID, treat as error */
my_error(ER_GALERA_REPLICATION_NOT_SUPPORTED, MYF(0));
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
ER_OPTION_PREVENTS_STATEMENT, "You can't execute statements that would generate local " "GTIDs when wsrep_mode = DISALLOW_LOCAL_GTID is set. " "Try disabling binary logging with SET sql_log_bin=0 " "to execute this statement.");
return -1;
}
if (thd->wsrep_cs().begin_rsu(5000))
{
WSREP_WARN("RSU begin failed");
} else
{
thd->variables.wsrep_on= 0;
} return0;
}
staticvoid wsrep_RSU_end(THD *thd)
{
WSREP_DEBUG("RSU END: %lld : %s", wsrep_thd_trx_seqno(thd),
wsrep_thd_query(thd)); if (thd->wsrep_cs().end_rsu())
{
WSREP_WARN("Failed to end RSU, server may need to be restarted");
}
thd->variables.wsrep_on= 1;
}
int wsrep_to_isolation_begin(THD *thd, constchar *db_, constchar *table_, const TABLE_LIST* table_list, const Alter_info *alter_info, const wsrep::key_array *fk_tables, const HA_CREATE_INFO *create_info)
{
DEBUG_SYNC(thd, "wsrep_kill_thd_before_enter_toi");
mysql_mutex_lock(&thd->LOCK_thd_kill); const killed_state killed = thd->killed;
mysql_mutex_unlock(&thd->LOCK_thd_kill); if (killed)
{ /* The thread may have been killed as a result of memory pressure. */ return -1;
}
/* Noisolationforapplierorreplayingthreads.
*/ if (!wsrep_thd_is_local(thd))
{ if (wsrep_OSU_method_get(thd) == WSREP_OSU_TOI)
WSREP_DEBUG("%s TOI Begin: %s",
is_replaying_connection(thd) ? "Replay" : "Apply",
wsrep_thd_query(thd));
return0;
}
if (thd->wsrep_parallel_slave_wait_for_prior_commit())
{
WSREP_WARN("TOI: wait_for_prior_commit() returned error."); return -1;
}
int ret= 0;
mysql_mutex_lock(&thd->LOCK_thd_data);
if (thd->wsrep_trx().state() == wsrep::transaction::s_must_abort)
{
WSREP_INFO("thread: %lld schema: %s query: %s has been aborted due to multi-master conflict",
(longlong) thd->thread_id, thd->get_db(), thd->query());
mysql_mutex_unlock(&thd->LOCK_thd_data); return WSREP_TRX_FAIL;
}
mysql_mutex_unlock(&thd->LOCK_thd_data);
if (Wsrep_server_state::instance().desynced_on_pause())
{
my_message(ER_UNKNOWN_COM_ERROR, "Aborting TOI: Replication paused on node for FTWRL/BACKUP STAGE.", MYF(0));
WSREP_DEBUG("Aborting TOI: Replication paused on node for FTWRL/BACKUP STAGE.: %s %llu",
wsrep_thd_query(thd), thd->thread_id); return -1;
}
/* If we are inside LOCK TABLE we release it and give warning. */ if (thd->variables.option_bits & OPTION_TABLE_LOCK &&
thd->lex->sql_command == SQLCOM_CREATE_SEQUENCE)
{
thd->locked_tables_list.unlock_locked_tables(thd);
thd->variables.option_bits&= ~(OPTION_TABLE_LOCK);
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN,
HA_ERR_UNSUPPORTED, "Galera cluster does not support LOCK TABLE on " "SEQUENCES. Lock is released.");
} if (wsrep_debug && thd->mdl_context.has_locks())
{
WSREP_DEBUG("thread holds MDL locks at TO begin: %s %llu",
wsrep_thd_query(thd), thd->thread_id);
}
if (wsrep_thd_is_aborting(granted_thd))
{ // Granted thread is aborting, we wait it to release MDL-locs
} elseif (wsrep_thd_is_SR(granted_thd) && wsrep_thd_is_toi(request_thd))
{ // Granted thread is executing streaming replication and request is DDL, // abort granted
wsrep_mdl_log(WSREP_MDL_INFO, "MDL conflict, DDL vs SR",
request_thd, granted_thd, ticket, key);
wsrep_abort_thd(request_thd, granted_thd, 1);
} else
{ // This case BF-BF is not possible so fail on debug
wsrep_mdl_log(WSREP_MDL_ERROR, "MDL BF-BF conflict",
request_thd, granted_thd, ticket, key);
DBUG_ASSERT(!(wsrep_thd_is_BF(granted_thd, false) &&
wsrep_thd_is_BF(request_thd, false)));
}
}
/** This function handles MDL-conflict when thread holding MDL-lock (granted_thd)hasongoing BACKUP ORFLUSHTABLESWITHREADLOCK ORFLUSHTABLESFOREXPORT ORLOCKTABLES
if (granted_thd->current_backup_stage != BACKUP_FINISHED &&
wsrep_check_mode(WSREP_MODE_BF_MARIABACKUP))
{ // User has allowed killing mariabackup
wsrep_mdl_log(WSREP_MDL_INFO, "MDL conflict, ongoing backup",
request_thd, granted_thd, ticket, key);
wsrep_abort_thd(request_thd, granted_thd, true);
} // else requestor must wait for MDL-lock to be released
}
/** This function handles MDL-conflict when thread holding MDL-lock (granted_thd)isnotBF(bruteforce)andrequestoris BF.
if (granted_thd->wsrep_trx().active())
{ // Granted thread has active wsrep transaction, abort it
wsrep_abort_thd(request_thd, granted_thd, 1);
} else
{ // Granted thread is not wsrep transaction or it has no // active transaction e.g. CREATE TABLE X AS SELECT, // signal KILL_QUERY and abort transaction
granted_thd->awake_no_mutex(KILL_QUERY_HARD);
ha_abort_transaction(request_thd, granted_thd, TRUE);
}
}
// Thread should not hold any mutexes below debug sync point
DEBUG_SYNC(request_thd, "before_wsrep_thd_abort");
DBUG_EXECUTE_IF("sync.before_wsrep_thd_abort", { constchar act[]= "now " "SIGNAL sync.before_wsrep_thd_abort_reached " "WAIT_FOR signal.before_wsrep_thd_abort";
DBUG_ASSERT(!debug_sync_set_action(request_thd, STRING_WITH_LEN(act)));
};);
/* We should hold THD::LOCK_thd_datatoprotectgrantedfromconcurrentusage andTHD::LOCK_thd_killtoprotectitfromdisconnectordelete.
*/
mysql_mutex_lock(&granted_thd->LOCK_thd_kill);
mysql_mutex_lock(&granted_thd->LOCK_thd_data);
if (granted_thd->wsrep_aborter != 0)
{ // Granted thread has being already selected as a victim for // BF kill, we can wait until it releases MDL-lock
DBUG_ASSERT(granted_thd->wsrep_aborter == request_thd->thread_id);
} elseif (granted_thd->current_backup_stage != BACKUP_FINISHED ||
granted_thd->global_read_lock.is_acquired() ||
granted_thd->locked_tables_mode == LTM_LOCK_TABLES ||
granted_thd->variables.option_bits & OPTION_TABLE_LOCK)
{ /* Granted thread has ongoingBACKUPOR ongoingFLUSHTABLESWITHREADLOCKOR ongoingFLUSHTABLESFOREXPORTOR ongoingLOCKTABLES
*/
wsrep_handle_locked(request_thd, granted_thd, ticket, key);
} elseif (wsrep_thd_is_BF(granted_thd, FALSE))
{ // Granted thread is either TOI or applying
wsrep_handle_granted_bf(request_thd, granted_thd, ticket, key);
} else
{ // Granted thread is not brute-force thread, we can abort it
wsrep_abort_granted(request_thd, granted_thd, ticket, key);
}
int wsrep_wait_committing_connections_close(int wait_time)
{ int sleep_time= 100;
WSREP_DEBUG("wait for committing transaction to close: %d sleep: %d", wait_time, sleep_time); while (server_threads.iterate(have_committing_connections) && wait_time > 0)
{
WSREP_DEBUG("wait for committing transaction to close: %d", wait_time);
my_sleep(sleep_time);
wait_time -= sleep_time;
} return server_threads.iterate(have_committing_connections);
}
static my_bool kill_all_threads(THD *thd, THD *caller_thd)
{
DBUG_PRINT("quit", ("Informing thread %lld that it's time to die",
(longlong) thd->thread_id)); /* We skip slave threads & scheduler on this first loop through. */ if (is_client_connection(thd) && thd != caller_thd)
{ /* the connection executing SHUTDOWN, should do clean exit,
not aborting here */ if (thd->get_command() == COM_SHUTDOWN)
{
WSREP_DEBUG("leaving SHUTDOWN executing connection alive, thread: %lld",
(longlong) thd->thread_id); return0;
} /* replaying connection is killed by signal */ if (is_replaying_connection(thd))
{
WSREP_DEBUG("closing connection is replaying %lld", (longlong) thd->thread_id);
thd->set_killed(KILL_CONNECTION_HARD); return0;
}
if (!abort_replicated(thd))
{ /* replicated transactions must be skipped */
WSREP_DEBUG("closing connection %lld", (longlong) thd->thread_id); /* instead of wsrep_close_thread() we do now hard kill by THD::awake */
thd->awake(KILL_CONNECTION_HARD); return0;
}
} return0;
}
DBUG_PRINT("quit", ("Waiting for threads to die (count=%u)",
THD_count::value()));
WSREP_DEBUG("waiting for client connections to close: %u",
THD_count::value());
while (wait_to_end && server_threads.iterate(have_client_connections))
{
sleep(1);
DBUG_PRINT("quit",("One thread died (count=%u)", THD_count::value()));
}
/* All client connection threads have now been aborted */
}
void wsrep_wait_appliers_close(THD *thd)
{ /* Wait for wsrep appliers to gracefully exit */
mysql_mutex_lock(&LOCK_wsrep_slave_threads); while (wsrep_running_threads > 2) /* 2isforrollbackerthreadwhichneedstobekilledexplicitly. Thisgottabefixedinamoreelegantmannerifwegonnahavearbitrary numberofnon-applierwsrepthreads.
*/
{
mysql_cond_wait(&COND_wsrep_slave_threads, &LOCK_wsrep_slave_threads);
}
mysql_mutex_unlock(&LOCK_wsrep_slave_threads);
DBUG_PRINT("quit",("applier threads have died (count=%u)",
uint32_t(wsrep_running_threads)));
/* Now kill remaining wsrep threads: rollbacker */
wsrep_close_threads (thd); /* and wait for them to die */
mysql_mutex_lock(&LOCK_wsrep_slave_threads); while (wsrep_running_threads > 0)
{
mysql_cond_wait(&COND_wsrep_slave_threads, &LOCK_wsrep_slave_threads);
}
mysql_mutex_unlock(&LOCK_wsrep_slave_threads);
DBUG_PRINT("quit",("all wsrep system threads have died"));
/* All wsrep applier threads have now been aborted. However, if this thread isalsoapplier,wearestillrunning...
*/
} int wsrep_must_ignore_error(THD* thd)
{ constint error= thd->get_stmt_da()->sql_errno(); const uint flags= sql_command_flags[thd->lex->sql_command];
DBUG_EXECUTE_IF("wsrep_simulate_failed_connection_1", goto error; ); // </5.1.17> /* handle_one_connection()isnormallytheonlywayathreadwould startandwouldalwaysbeontheveryhighendofthestack, therefore,thethreadstackalwaysstartsattheaddressofthe firstlocalvariableofhandle_one_connection,whichisthd.We needtoknowthestartofthestacksothatwecouldcheckfor stackoverruns.
*/
DBUG_PRINT("wsrep", ("handle_one_connection called by thread %lld",
(longlong)thd->thread_id)); /* now that we've called my_thread_init(), it is safe to call DBUG_* */
error:
WSREP_ERROR("Failed to create/initialize system thread");
if (thd)
{
close_connection(thd, ER_OUT_OF_RESOURCES);
statistic_increment(aborted_connects, &LOCK_status);
server_threads.erase(thd); delete thd;
my_thread_end();
} delete thd_args; // This will signal error to wsrep_slave_threads_update
wsrep_thread_create_failed.store(true, std::memory_order_relaxed);
/* Abort if its the first applier/rollbacker thread. */ if (!mysqld_server_initialized)
unireg_abort(1); else return NULL;
}
enum wsrep::streaming_context::fragment_unit wsrep_fragment_unit(ulong unit)
{ switch (unit)
{ case WSREP_FRAG_BYTES: return wsrep::streaming_context::bytes; case WSREP_FRAG_ROWS: return wsrep::streaming_context::row; case WSREP_FRAG_STATEMENTS: return wsrep::streaming_context::statement; default:
DBUG_ASSERT(0); return wsrep::streaming_context::bytes;
}
}
bool wsrep_wait_ready(THD *thd)
{ // First check not locking the mutex. switch (wsrep_state.load()) { case NOT_READY: break; case READY: returntrue; case IN_SHUTDOWN: returnfalse;
}
mysql_mutex_lock(&LOCK_wsrep_ready); while(!wsrep_state)
{
WSREP_INFO("Waiting to reach ready state");
mysql_cond_wait(&COND_wsrep_ready, &LOCK_wsrep_ready);
}
mysql_mutex_unlock(&LOCK_wsrep_ready);
WSREP_INFO("ready state reached"); /* It may happen we stop waiting and immediately transition tonotreadystatewhenwereachthisline.It'saspurious wakeupwecannotdealwith,butthebestwecandoistoreturn readinessindication.That'swhyweshouldcompareto
IN_SHUTDOWN rather than returning wsrep_ready_get(). */ return wsrep_state != IN_SHUTDOWN;
}
void wsrep_ready_set(bool ready_value)
{
WSREP_DEBUG("Setting wsrep_ready to %d", ready_value);
mysql_mutex_lock(&LOCK_wsrep_ready); /* Only transition if we're not shutting down. */ if (wsrep_state != IN_SHUTDOWN) {
wsrep_state= ready_value ? READY : NOT_READY; // Signal if we have reached ready state if (ready_value)
mysql_cond_signal(&COND_wsrep_ready);
}
mysql_mutex_unlock(&LOCK_wsrep_ready);
}
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.