Loading sql/log_event.cc +61 −29 Original line number Diff line number Diff line Loading @@ -22,11 +22,11 @@ #include "mysql_priv.h" #endif /* MYSQL_CLIENT */ #define LOG_EVENT_HEADER_LEN 9 #define LOG_EVENT_HEADER_LEN 13 #define QUERY_HEADER_LEN (sizeof(uint32) + sizeof(uint32) + sizeof(uchar)) #define LOAD_HEADER_LEN (sizeof(uint32) + sizeof(uint32) + \ + sizeof(uint32) + 2 + sizeof(uint32)) #define EVENT_LEN_OFFSET 5 #define EVENT_LEN_OFFSET 9 #define EVENT_TYPE_OFFSET 4 #define MAX_EVENT_LEN 4*1024*1024 #define QUERY_EVENT_OVERHEAD LOG_EVENT_HEADER_LEN+QUERY_HEADER_LEN Loading Loading @@ -71,6 +71,8 @@ int Log_event::write_header(FILE* file) int4store(pos, when); // timestamp pos += 4; *pos++ = get_type_code(); // event type code int4store(pos, server_id); pos += 4; int4store(pos, get_data_size() + LOG_EVENT_HEADER_LEN); pos += 4; return (my_fwrite(file, (byte*) buf, (uint) (pos - buf), Loading Loading @@ -106,16 +108,19 @@ int Log_event::read_log_event(FILE* file, String* packet) Log_event* Log_event::read_log_event(FILE* file) { time_t timestamp; char buf[5]; uint32 server_id; char buf[LOG_EVENT_HEADER_LEN-4]; if (my_fread(file, (byte *) buf, sizeof(buf), MY_NABP)) return NULL; timestamp = uint4korr(buf); server_id = uint4korr(buf + 5); switch(buf[EVENT_TYPE_OFFSET]) { case QUERY_EVENT: { Query_log_event* q = new Query_log_event(file, timestamp); Query_log_event* q = new Query_log_event(file, timestamp, server_id); if (!q->query) { delete q; Loading @@ -127,7 +132,7 @@ Log_event* Log_event::read_log_event(FILE* file) case LOAD_EVENT: { Load_log_event* l = new Load_log_event(file, timestamp); Load_log_event* l = new Load_log_event(file, timestamp, server_id); if (!l->table_name) { delete l; Loading @@ -140,7 +145,7 @@ Log_event* Log_event::read_log_event(FILE* file) case ROTATE_EVENT: { Rotate_log_event* r = new Rotate_log_event(file, timestamp); Rotate_log_event* r = new Rotate_log_event(file, timestamp, server_id); if (!r->new_log_ident) { delete r; Loading @@ -152,7 +157,7 @@ Log_event* Log_event::read_log_event(FILE* file) case INTVAR_EVENT: { Intvar_log_event* e = new Intvar_log_event(file, timestamp); Intvar_log_event* e = new Intvar_log_event(file, timestamp, server_id); if (e->type == INVALID_INT_EVENT) { delete e; Loading @@ -162,8 +167,8 @@ Log_event* Log_event::read_log_event(FILE* file) return e; } case START_EVENT: return new Start_log_event(file, timestamp); case STOP_EVENT: return new Stop_log_event(file, timestamp); case START_EVENT: return new Start_log_event(file, timestamp, server_id); case STOP_EVENT: return new Stop_log_event(file, timestamp, server_id); default: return NULL; } Loading Loading @@ -221,12 +226,23 @@ Log_event* Log_event::read_log_event(const char* buf, int max_buf) return NULL; } void Log_event::print_timestamp(FILE* file) void Log_event::print_header(FILE* file) { fputc('#', file); print_timestamp(file); fprintf(file, " server id %d ", server_id); } void Log_event::print_timestamp(FILE* file, time_t* ts = 0) { struct tm tm_tmp; localtime_r(&when,&tm_tmp); if(!ts) { ts = &when; } localtime_r(ts,&tm_tmp); fprintf(file,"#%02d%02d%02d %2d:%02d:%02d", fprintf(file,"%02d%02d%02d %2d:%02d:%02d", tm_tmp.tm_year % 100, tm_tmp.tm_mon+1, tm_tmp.tm_mday, Loading @@ -241,8 +257,11 @@ void Start_log_event::print(FILE* file, bool short_form) if (short_form) return; print_timestamp(file); fprintf(file, "\tStart\n"); print_header(file); fprintf(file, "\tStart: binlog v %d, server v %s created ", binlog_version, server_version); print_timestamp(file, (time_t*)&created); fputc('\n', file); fflush(file); } Loading @@ -251,7 +270,7 @@ void Stop_log_event::print(FILE* file, bool short_form) if (short_form) return; print_timestamp(file); print_header(file); fprintf(file, "\tStop\n"); fflush(file); } Loading @@ -261,7 +280,7 @@ void Rotate_log_event::print(FILE* file, bool short_form) if (short_form) return; print_timestamp(file); print_header(file); fprintf(file, "\tRotate to "); if (new_log_ident) my_fwrite(file, (byte*) new_log_ident, (uint)ident_len, Loading @@ -270,8 +289,9 @@ void Rotate_log_event::print(FILE* file, bool short_form) fflush(file); } Rotate_log_event::Rotate_log_event(FILE* file, time_t when_arg): Log_event(when_arg),new_log_ident(NULL),alloced(0) Rotate_log_event::Rotate_log_event(FILE* file, time_t when_arg, uint32 server_id): Log_event(when_arg, 0, 0, server_id),new_log_ident(NULL),alloced(0) { char *tmp_ident; char buf[4]; Loading @@ -298,6 +318,14 @@ Rotate_log_event::Rotate_log_event(FILE* file, time_t when_arg): alloced = 1; } Start_log_event::Start_log_event(const char* buf) :Log_event(buf) { buf += EVENT_LEN_OFFSET + 4; // skip even length binlog_version = uint2korr(buf); memcpy(server_version, buf + 2, sizeof(server_version)); created = uint4korr(buf + 2 + sizeof(server_version)); } Rotate_log_event::Rotate_log_event(const char* buf, int max_buf): Log_event(buf),new_log_ident(NULL),alloced(0) { Loading @@ -322,8 +350,9 @@ int Rotate_log_event::write_data(FILE* file) return 0; } Query_log_event::Query_log_event(FILE* file, time_t when_arg): Log_event(when_arg),data_buf(0),query(NULL),db(NULL) Query_log_event::Query_log_event(FILE* file, time_t when_arg, uint32 server_id): Log_event(when_arg,0,0,server_id),data_buf(0),query(NULL),db(NULL) { char buf[QUERY_HEADER_LEN + 4]; ulong data_len; Loading Loading @@ -382,7 +411,7 @@ void Query_log_event::print(FILE* file, bool short_form) { if (!short_form) { print_timestamp(file); print_header(file); fprintf(file, "\tQuery\tthread_id=%lu\texec_time=%lu\n", (ulong) thread_id, (ulong) exec_time); } Loading Loading @@ -414,8 +443,9 @@ int Query_log_event::write_data(FILE* file) return 0; } Intvar_log_event:: Intvar_log_event(FILE* file, time_t when_arg) :Log_event(when_arg), type(INVALID_INT_EVENT) Intvar_log_event:: Intvar_log_event(FILE* file, time_t when_arg, uint32 server_id) :Log_event(when_arg,0,0,server_id), type(INVALID_INT_EVENT) { my_fseek(file, 4L, MY_SEEK_CUR, MYF(MY_WME)); // skip the event length char buf[9]; Loading Loading @@ -444,7 +474,7 @@ void Intvar_log_event::print(FILE* file, bool short_form) char llbuff[22]; if(!short_form) { print_timestamp(file); print_header(file); fprintf(file, "\tIntvar\n"); } Loading Loading @@ -493,8 +523,9 @@ int Load_log_event::write_data(FILE* file __attribute__((unused))) return 0; } Load_log_event::Load_log_event(FILE* file, time_t when): Log_event(when),data_buf(0),num_fields(0),fields(0),field_lens(0),field_block_len(0), Load_log_event::Load_log_event(FILE* file, time_t when, uint32 server_id): Log_event(when,0,0,server_id),data_buf(0),num_fields(0), fields(0),field_lens(0),field_block_len(0), table_name(0),db(0),fname(0) { Loading Loading @@ -539,7 +570,8 @@ Load_log_event::Load_log_event(FILE* file, time_t when): } Load_log_event::Load_log_event(const char* buf, int max_buf): Log_event(when),data_buf(0),num_fields(0),fields(0),field_lens(0),field_block_len(0), Log_event(when,0,0,server_id),data_buf(0),num_fields(0),fields(0), field_lens(0),field_block_len(0), table_name(0),db(0),fname(0) { Loading Loading @@ -594,7 +626,7 @@ void Load_log_event::print(FILE* file, bool short_form) { if (!short_form) { print_timestamp(file); print_header(file); fprintf(file, "\tQuery\tthread_id=%d\texec_time=%ld\n", thread_id, exec_time); } Loading sql/log_event.h +58 −18 Original line number Diff line number Diff line Loading @@ -28,7 +28,7 @@ #define LOG_READ_MEM -5 #define LOG_EVENT_OFFSET 4 #define BINLOG_VERSION 1 enum Log_event_type { START_EVENT = 1, QUERY_EVENT =2, STOP_EVENT=3, ROTATE_EVENT = 4, INTVAR_EVENT=5, Loading @@ -40,25 +40,32 @@ enum Int_event_type { INVALID_INT_EVENT = 0, LAST_INSERT_ID_EVENT = 1, INSERT_ID class String; #endif extern uint32 server_id; class Log_event { public: time_t when; ulong exec_time; int valid_exec_time; // if false, the exec time setting is bogus and needs int valid_exec_time; // if false, the exec time setting is bogus uint32 server_id; int write(FILE* file); int write_header(FILE* file); virtual int write_data(FILE* file __attribute__((unused))) { return 0; } virtual Log_event_type get_type_code() = 0; Log_event(time_t when_arg, ulong exec_time_arg = 0, int valid_exec_time_arg = 0): when(when_arg), exec_time(exec_time_arg), valid_exec_time(valid_exec_time_arg) {} int valid_exec_time_arg = 0, uint32 server_id = 0): when(when_arg), exec_time(exec_time_arg), valid_exec_time(valid_exec_time_arg) { if(!server_id) this->server_id = ::server_id; else this->server_id = server_id; } Log_event(const char* buf): valid_exec_time(1) Log_event(const char* buf): valid_exec_time(0) { when = uint4korr(buf); exec_time = uint4korr(buf + 5); server_id = uint4korr(buf + 5); } virtual ~Log_event() {} Loading @@ -66,7 +73,8 @@ class Log_event virtual int get_data_size() { return 0;} virtual void print(FILE* file, bool short_form = 0) = 0; void print_timestamp(FILE* file); void print_timestamp(FILE* file, time_t *ts = 0); void print_header(FILE* file); static Log_event* read_log_event(FILE* file); static Log_event* read_log_event(const char* buf, int max_buf); Loading @@ -93,7 +101,7 @@ class Query_log_event: public Log_event #if !defined(MYSQL_CLIENT) THD* thd; Query_log_event(THD* thd_arg, const char* query_arg): Log_event(thd_arg->start_time), data_buf(0), Log_event(thd_arg->start_time,0,0,thd->server_id), data_buf(0), query(query_arg), db(thd_arg->db), q_len(thd_arg->query_length), thread_id(thd_arg->thread_id), thd(thd_arg) { Loading @@ -105,7 +113,7 @@ class Query_log_event: public Log_event } #endif Query_log_event(FILE* file, time_t when); Query_log_event(FILE* file, time_t when, uint32 server_id); Query_log_event(const char* buf, int max_buf); ~Query_log_event() { Loading Loading @@ -244,7 +252,7 @@ class Load_log_event: public Log_event void set_fields(List<Item> &fields); #endif Load_log_event(FILE* file, time_t when); Load_log_event(FILE* file, time_t when, uint32 server_id); Load_log_event(const char* buf, int max_buf); ~Load_log_event() { Loading @@ -269,21 +277,52 @@ class Load_log_event: public Log_event void print(FILE* file, bool short_form = 0); }; extern char server_version[50]; class Start_log_event: public Log_event { public: Start_log_event() :Log_event(time(NULL)) {} Start_log_event(FILE* file, time_t when_arg) :Log_event(when_arg) uint16 binlog_version; char server_version[50]; uint32 created; Start_log_event() :Log_event(time(NULL)),binlog_version(BINLOG_VERSION) { my_fseek(file, 4L, MY_SEEK_CUR, MYF(MY_WME)); // skip the event length created = when; memcpy(server_version, ::server_version, sizeof(server_version)); } Start_log_event(const char* buf) :Log_event(buf) Start_log_event(FILE* file, time_t when_arg, uint32 server_id) : Log_event(when_arg, 0, 0, server_id) { char buf[sizeof(server_version) + sizeof(binlog_version) + sizeof(created)]; my_fseek(file, 4L, MY_SEEK_CUR, MYF(MY_WME)); // skip the event length if (my_fread(file, (byte*) buf, sizeof(buf), MYF(MY_NABP | MY_WME))) return; binlog_version = uint2korr(buf); memcpy(server_version, buf + 2, sizeof(server_version)); created = uint4korr(buf + 2 + sizeof(server_version)); } Start_log_event(const char* buf); ~Start_log_event() {} Log_event_type get_type_code() { return START_EVENT;} int write_data(FILE* file) { if(my_fwrite(file, (byte*) &binlog_version, sizeof(binlog_version), MYF(MY_NABP | MY_WME)) || my_fwrite(file, (byte*) server_version, sizeof(server_version), MYF(MY_NABP | MY_WME)) || my_fwrite(file, (byte*) &created, sizeof(created), MYF(MY_NABP | MY_WME))) return -1; return 0; } int get_data_size() { return sizeof(binlog_version) + sizeof(server_version) + sizeof(created); } void print(FILE* file, bool short_form = 0); }; Loading @@ -295,7 +334,7 @@ class Intvar_log_event: public Log_event Intvar_log_event(uchar type_arg, ulonglong val_arg) :Log_event(time(NULL)),val(val_arg),type(type_arg) {} Intvar_log_event(FILE* file, time_t when); Intvar_log_event(FILE* file, time_t when, uint32 server_id); Intvar_log_event(const char* buf); ~Intvar_log_event() {} Log_event_type get_type_code() { return INTVAR_EVENT;} Loading @@ -311,7 +350,8 @@ class Stop_log_event: public Log_event public: Stop_log_event() :Log_event(time(NULL)) {} Stop_log_event(FILE* file, time_t when_arg): Log_event(when_arg) Stop_log_event(FILE* file, time_t when_arg, uint32 server_id): Log_event(when_arg,0,0,server_id) { my_fseek(file, 4L, MY_SEEK_CUR, MYF(MY_WME)); // skip the event length } Loading @@ -337,7 +377,7 @@ class Rotate_log_event: public Log_event alloced(0) {} Rotate_log_event(FILE* file, time_t when) ; Rotate_log_event(FILE* file, time_t when, uint32 server_id) ; Rotate_log_event(const char* buf, int max_buf); ~Rotate_log_event() { Loading sql/mysql_priv.h +1 −0 Original line number Diff line number Diff line Loading @@ -437,6 +437,7 @@ void sql_perror(const char *message); void sql_print_error(const char *format,...) __attribute__ ((format (printf, 1, 2))); extern uint32 server_id; extern char mysql_data_home[2],server_version[50],max_sort_char, mysql_real_data_home[]; extern my_string mysql_unix_port,mysql_tmpdir; Loading sql/mysqlbinlog.cc +2 −0 Original line number Diff line number Diff line Loading @@ -29,6 +29,8 @@ #define CLIENT_CAPABILITIES (CLIENT_LONG_PASSWORD | CLIENT_LONG_FLAG | CLIENT_LOCAL_FILES) char server_version[50]; uint32 server_id = 0; // needed by net_serv.c ulong bytes_sent = 0L, bytes_received = 0L; Loading sql/mysqld.cc +39 −14 Original line number Diff line number Diff line Loading @@ -142,6 +142,7 @@ static uint handler_count; static bool opt_console=0; #endif static bool opt_skip_slave_start = 0; // if set, slave is not autostarted static ulong opt_specialflag=SPECIAL_ENGLISH; static my_socket unix_sock= INVALID_SOCKET,ip_sock= INVALID_SOCKET; static ulong back_log,connect_timeout,concurrency; Loading Loading @@ -180,6 +181,7 @@ I_List<i_string> replicate_do_db, replicate_ignore_db; // allow the user to tell us which db to replicate and which to ignore I_List<i_string> binlog_do_db, binlog_ignore_db; uint32 server_id = 0; // server id for replication uint mysql_port; uint test_flags, select_errors=0, dropping_tables=0,ha_open_options=0; uint volatile thread_count=0, thread_running=0, kill_cached_threads=0, Loading Loading @@ -1480,6 +1482,8 @@ int main(int argc, char **argv) open_log(&mysql_update_log, hostname, opt_update_logname, "", LOG_NEW); if (opt_bin_log) { if(server_id) { if (!opt_bin_logname) { Loading @@ -1492,6 +1496,10 @@ int main(int argc, char **argv) open_log(&mysql_bin_log, hostname, opt_bin_logname, "-bin", LOG_BIN); } else sql_print_error("Server id is not set - binary logging disabled"); } if (opt_slow_log) open_log(&mysql_slow_log, hostname, opt_slow_logname, "-slow.log", LOG_NORMAL); Loading Loading @@ -1594,11 +1602,16 @@ int main(int argc, char **argv) // slave thread if(master_host) { if(server_id) { pthread_t hThread; if(pthread_create(&hThread, &connection_attrib, handle_slave, 0)) if(!opt_skip_slave_start && pthread_create(&hThread, &connection_attrib, handle_slave, 0)) sql_print_error("Warning: Can't create thread to handle slave"); } else sql_print_error("Server id is not set, slave thread will not be started"); } printf(ER(ER_READY),my_progname,server_version,""); Loading Loading @@ -2201,6 +2214,7 @@ enum options { OPT_LOG_SLAVE_UPDATES, OPT_BINLOG_DO_DB, OPT_BINLOG_IGNORE_DB, OPT_WANT_CORE, OPT_SKIP_CONCURRENT_INSERT, OPT_MEMLOCK, OPT_MYISAM_RECOVER, OPT_REPLICATE_REWRITE_DB, OPT_SERVER_ID, OPT_SKIP_SLAVE_START }; static struct option long_options[] = { Loading Loading @@ -2265,8 +2279,11 @@ static struct option long_options[] = { {"port", required_argument, 0, 'P'}, {"replicate-do-db", required_argument, 0, (int) OPT_REPLICATE_DO_DB}, {"replicate-ignore-db", required_argument, 0, (int) OPT_REPLICATE_IGNORE_DB}, {"replicate-rewrite-db", required_argument, 0, (int) OPT_REPLICATE_REWRITE_DB}, {"safe-mode", no_argument, 0, (int) OPT_SAFE}, {"socket", required_argument, 0, (int) OPT_SOCKET}, {"server-id", required_argument, 0, (int)OPT_SERVER_ID}, {"set-variable", required_argument, 0, 'O'}, #ifdef HAVE_BERKELEY_DB {"skip-bdb", no_argument, 0, (int) OPT_BDB_SKIP}, Loading @@ -2279,6 +2296,7 @@ static struct option long_options[] = { {"skip-name-resolve", no_argument, 0, (int) OPT_SKIP_RESOLVE}, {"skip-new", no_argument, 0, (int) OPT_SKIP_NEW}, {"skip-show-database", no_argument, 0, (int) OPT_SKIP_SHOW_DB}, {"skip-slave-start", no_argument, 0, (int) OPT_SKIP_SLAVE_START}, {"skip-networking", no_argument, 0, (int) OPT_SKIP_NETWORKING}, {"skip-thread-priority", no_argument, 0, (int) OPT_SKIP_PRIOR}, {"sql-bin-update-same", no_argument, 0, (int) OPT_SQL_BIN_UPDATE_SAME}, Loading Loading @@ -2427,6 +2445,7 @@ struct show_var_st init_vars[]= { {"port", (char*) &mysql_port, SHOW_INT}, {"protocol_version", (char*) &protocol_version, SHOW_INT}, {"record_buffer", (char*) &my_default_record_cache_size,SHOW_LONG}, {"server_id", (char*) &server_id, SHOW_LONG}, {"skip_locking", (char*) &my_disable_locking, SHOW_MY_BOOL}, {"skip_networking", (char*) &opt_disable_networking, SHOW_BOOL}, {"skip_show_database", (char*) &opt_skip_show_db, SHOW_BOOL}, Loading Loading @@ -2850,6 +2869,9 @@ static void get_options(int argc,char **argv) opt_slow_log=1; opt_slow_logname=optarg; break; case (int)OPT_SKIP_SLAVE_START: opt_skip_slave_start = 1; break; case (int) OPT_SKIP_NEW: opt_specialflag|= SPECIAL_NO_NEW_FUNC; default_table_type=DB_TYPE_ISAM; Loading Loading @@ -2971,6 +2993,9 @@ static void get_options(int argc,char **argv) default_table_type= (enum db_type) type; break; } case OPT_SERVER_ID: server_id = atoi(optarg); break; case OPT_DELAY_KEY_WRITE: ha_open_options|=HA_OPEN_DELAY_KEY_WRITE; myisam_delay_key_write=1; Loading Loading
sql/log_event.cc +61 −29 Original line number Diff line number Diff line Loading @@ -22,11 +22,11 @@ #include "mysql_priv.h" #endif /* MYSQL_CLIENT */ #define LOG_EVENT_HEADER_LEN 9 #define LOG_EVENT_HEADER_LEN 13 #define QUERY_HEADER_LEN (sizeof(uint32) + sizeof(uint32) + sizeof(uchar)) #define LOAD_HEADER_LEN (sizeof(uint32) + sizeof(uint32) + \ + sizeof(uint32) + 2 + sizeof(uint32)) #define EVENT_LEN_OFFSET 5 #define EVENT_LEN_OFFSET 9 #define EVENT_TYPE_OFFSET 4 #define MAX_EVENT_LEN 4*1024*1024 #define QUERY_EVENT_OVERHEAD LOG_EVENT_HEADER_LEN+QUERY_HEADER_LEN Loading Loading @@ -71,6 +71,8 @@ int Log_event::write_header(FILE* file) int4store(pos, when); // timestamp pos += 4; *pos++ = get_type_code(); // event type code int4store(pos, server_id); pos += 4; int4store(pos, get_data_size() + LOG_EVENT_HEADER_LEN); pos += 4; return (my_fwrite(file, (byte*) buf, (uint) (pos - buf), Loading Loading @@ -106,16 +108,19 @@ int Log_event::read_log_event(FILE* file, String* packet) Log_event* Log_event::read_log_event(FILE* file) { time_t timestamp; char buf[5]; uint32 server_id; char buf[LOG_EVENT_HEADER_LEN-4]; if (my_fread(file, (byte *) buf, sizeof(buf), MY_NABP)) return NULL; timestamp = uint4korr(buf); server_id = uint4korr(buf + 5); switch(buf[EVENT_TYPE_OFFSET]) { case QUERY_EVENT: { Query_log_event* q = new Query_log_event(file, timestamp); Query_log_event* q = new Query_log_event(file, timestamp, server_id); if (!q->query) { delete q; Loading @@ -127,7 +132,7 @@ Log_event* Log_event::read_log_event(FILE* file) case LOAD_EVENT: { Load_log_event* l = new Load_log_event(file, timestamp); Load_log_event* l = new Load_log_event(file, timestamp, server_id); if (!l->table_name) { delete l; Loading @@ -140,7 +145,7 @@ Log_event* Log_event::read_log_event(FILE* file) case ROTATE_EVENT: { Rotate_log_event* r = new Rotate_log_event(file, timestamp); Rotate_log_event* r = new Rotate_log_event(file, timestamp, server_id); if (!r->new_log_ident) { delete r; Loading @@ -152,7 +157,7 @@ Log_event* Log_event::read_log_event(FILE* file) case INTVAR_EVENT: { Intvar_log_event* e = new Intvar_log_event(file, timestamp); Intvar_log_event* e = new Intvar_log_event(file, timestamp, server_id); if (e->type == INVALID_INT_EVENT) { delete e; Loading @@ -162,8 +167,8 @@ Log_event* Log_event::read_log_event(FILE* file) return e; } case START_EVENT: return new Start_log_event(file, timestamp); case STOP_EVENT: return new Stop_log_event(file, timestamp); case START_EVENT: return new Start_log_event(file, timestamp, server_id); case STOP_EVENT: return new Stop_log_event(file, timestamp, server_id); default: return NULL; } Loading Loading @@ -221,12 +226,23 @@ Log_event* Log_event::read_log_event(const char* buf, int max_buf) return NULL; } void Log_event::print_timestamp(FILE* file) void Log_event::print_header(FILE* file) { fputc('#', file); print_timestamp(file); fprintf(file, " server id %d ", server_id); } void Log_event::print_timestamp(FILE* file, time_t* ts = 0) { struct tm tm_tmp; localtime_r(&when,&tm_tmp); if(!ts) { ts = &when; } localtime_r(ts,&tm_tmp); fprintf(file,"#%02d%02d%02d %2d:%02d:%02d", fprintf(file,"%02d%02d%02d %2d:%02d:%02d", tm_tmp.tm_year % 100, tm_tmp.tm_mon+1, tm_tmp.tm_mday, Loading @@ -241,8 +257,11 @@ void Start_log_event::print(FILE* file, bool short_form) if (short_form) return; print_timestamp(file); fprintf(file, "\tStart\n"); print_header(file); fprintf(file, "\tStart: binlog v %d, server v %s created ", binlog_version, server_version); print_timestamp(file, (time_t*)&created); fputc('\n', file); fflush(file); } Loading @@ -251,7 +270,7 @@ void Stop_log_event::print(FILE* file, bool short_form) if (short_form) return; print_timestamp(file); print_header(file); fprintf(file, "\tStop\n"); fflush(file); } Loading @@ -261,7 +280,7 @@ void Rotate_log_event::print(FILE* file, bool short_form) if (short_form) return; print_timestamp(file); print_header(file); fprintf(file, "\tRotate to "); if (new_log_ident) my_fwrite(file, (byte*) new_log_ident, (uint)ident_len, Loading @@ -270,8 +289,9 @@ void Rotate_log_event::print(FILE* file, bool short_form) fflush(file); } Rotate_log_event::Rotate_log_event(FILE* file, time_t when_arg): Log_event(when_arg),new_log_ident(NULL),alloced(0) Rotate_log_event::Rotate_log_event(FILE* file, time_t when_arg, uint32 server_id): Log_event(when_arg, 0, 0, server_id),new_log_ident(NULL),alloced(0) { char *tmp_ident; char buf[4]; Loading @@ -298,6 +318,14 @@ Rotate_log_event::Rotate_log_event(FILE* file, time_t when_arg): alloced = 1; } Start_log_event::Start_log_event(const char* buf) :Log_event(buf) { buf += EVENT_LEN_OFFSET + 4; // skip even length binlog_version = uint2korr(buf); memcpy(server_version, buf + 2, sizeof(server_version)); created = uint4korr(buf + 2 + sizeof(server_version)); } Rotate_log_event::Rotate_log_event(const char* buf, int max_buf): Log_event(buf),new_log_ident(NULL),alloced(0) { Loading @@ -322,8 +350,9 @@ int Rotate_log_event::write_data(FILE* file) return 0; } Query_log_event::Query_log_event(FILE* file, time_t when_arg): Log_event(when_arg),data_buf(0),query(NULL),db(NULL) Query_log_event::Query_log_event(FILE* file, time_t when_arg, uint32 server_id): Log_event(when_arg,0,0,server_id),data_buf(0),query(NULL),db(NULL) { char buf[QUERY_HEADER_LEN + 4]; ulong data_len; Loading Loading @@ -382,7 +411,7 @@ void Query_log_event::print(FILE* file, bool short_form) { if (!short_form) { print_timestamp(file); print_header(file); fprintf(file, "\tQuery\tthread_id=%lu\texec_time=%lu\n", (ulong) thread_id, (ulong) exec_time); } Loading Loading @@ -414,8 +443,9 @@ int Query_log_event::write_data(FILE* file) return 0; } Intvar_log_event:: Intvar_log_event(FILE* file, time_t when_arg) :Log_event(when_arg), type(INVALID_INT_EVENT) Intvar_log_event:: Intvar_log_event(FILE* file, time_t when_arg, uint32 server_id) :Log_event(when_arg,0,0,server_id), type(INVALID_INT_EVENT) { my_fseek(file, 4L, MY_SEEK_CUR, MYF(MY_WME)); // skip the event length char buf[9]; Loading Loading @@ -444,7 +474,7 @@ void Intvar_log_event::print(FILE* file, bool short_form) char llbuff[22]; if(!short_form) { print_timestamp(file); print_header(file); fprintf(file, "\tIntvar\n"); } Loading Loading @@ -493,8 +523,9 @@ int Load_log_event::write_data(FILE* file __attribute__((unused))) return 0; } Load_log_event::Load_log_event(FILE* file, time_t when): Log_event(when),data_buf(0),num_fields(0),fields(0),field_lens(0),field_block_len(0), Load_log_event::Load_log_event(FILE* file, time_t when, uint32 server_id): Log_event(when,0,0,server_id),data_buf(0),num_fields(0), fields(0),field_lens(0),field_block_len(0), table_name(0),db(0),fname(0) { Loading Loading @@ -539,7 +570,8 @@ Load_log_event::Load_log_event(FILE* file, time_t when): } Load_log_event::Load_log_event(const char* buf, int max_buf): Log_event(when),data_buf(0),num_fields(0),fields(0),field_lens(0),field_block_len(0), Log_event(when,0,0,server_id),data_buf(0),num_fields(0),fields(0), field_lens(0),field_block_len(0), table_name(0),db(0),fname(0) { Loading Loading @@ -594,7 +626,7 @@ void Load_log_event::print(FILE* file, bool short_form) { if (!short_form) { print_timestamp(file); print_header(file); fprintf(file, "\tQuery\tthread_id=%d\texec_time=%ld\n", thread_id, exec_time); } Loading
sql/log_event.h +58 −18 Original line number Diff line number Diff line Loading @@ -28,7 +28,7 @@ #define LOG_READ_MEM -5 #define LOG_EVENT_OFFSET 4 #define BINLOG_VERSION 1 enum Log_event_type { START_EVENT = 1, QUERY_EVENT =2, STOP_EVENT=3, ROTATE_EVENT = 4, INTVAR_EVENT=5, Loading @@ -40,25 +40,32 @@ enum Int_event_type { INVALID_INT_EVENT = 0, LAST_INSERT_ID_EVENT = 1, INSERT_ID class String; #endif extern uint32 server_id; class Log_event { public: time_t when; ulong exec_time; int valid_exec_time; // if false, the exec time setting is bogus and needs int valid_exec_time; // if false, the exec time setting is bogus uint32 server_id; int write(FILE* file); int write_header(FILE* file); virtual int write_data(FILE* file __attribute__((unused))) { return 0; } virtual Log_event_type get_type_code() = 0; Log_event(time_t when_arg, ulong exec_time_arg = 0, int valid_exec_time_arg = 0): when(when_arg), exec_time(exec_time_arg), valid_exec_time(valid_exec_time_arg) {} int valid_exec_time_arg = 0, uint32 server_id = 0): when(when_arg), exec_time(exec_time_arg), valid_exec_time(valid_exec_time_arg) { if(!server_id) this->server_id = ::server_id; else this->server_id = server_id; } Log_event(const char* buf): valid_exec_time(1) Log_event(const char* buf): valid_exec_time(0) { when = uint4korr(buf); exec_time = uint4korr(buf + 5); server_id = uint4korr(buf + 5); } virtual ~Log_event() {} Loading @@ -66,7 +73,8 @@ class Log_event virtual int get_data_size() { return 0;} virtual void print(FILE* file, bool short_form = 0) = 0; void print_timestamp(FILE* file); void print_timestamp(FILE* file, time_t *ts = 0); void print_header(FILE* file); static Log_event* read_log_event(FILE* file); static Log_event* read_log_event(const char* buf, int max_buf); Loading @@ -93,7 +101,7 @@ class Query_log_event: public Log_event #if !defined(MYSQL_CLIENT) THD* thd; Query_log_event(THD* thd_arg, const char* query_arg): Log_event(thd_arg->start_time), data_buf(0), Log_event(thd_arg->start_time,0,0,thd->server_id), data_buf(0), query(query_arg), db(thd_arg->db), q_len(thd_arg->query_length), thread_id(thd_arg->thread_id), thd(thd_arg) { Loading @@ -105,7 +113,7 @@ class Query_log_event: public Log_event } #endif Query_log_event(FILE* file, time_t when); Query_log_event(FILE* file, time_t when, uint32 server_id); Query_log_event(const char* buf, int max_buf); ~Query_log_event() { Loading Loading @@ -244,7 +252,7 @@ class Load_log_event: public Log_event void set_fields(List<Item> &fields); #endif Load_log_event(FILE* file, time_t when); Load_log_event(FILE* file, time_t when, uint32 server_id); Load_log_event(const char* buf, int max_buf); ~Load_log_event() { Loading @@ -269,21 +277,52 @@ class Load_log_event: public Log_event void print(FILE* file, bool short_form = 0); }; extern char server_version[50]; class Start_log_event: public Log_event { public: Start_log_event() :Log_event(time(NULL)) {} Start_log_event(FILE* file, time_t when_arg) :Log_event(when_arg) uint16 binlog_version; char server_version[50]; uint32 created; Start_log_event() :Log_event(time(NULL)),binlog_version(BINLOG_VERSION) { my_fseek(file, 4L, MY_SEEK_CUR, MYF(MY_WME)); // skip the event length created = when; memcpy(server_version, ::server_version, sizeof(server_version)); } Start_log_event(const char* buf) :Log_event(buf) Start_log_event(FILE* file, time_t when_arg, uint32 server_id) : Log_event(when_arg, 0, 0, server_id) { char buf[sizeof(server_version) + sizeof(binlog_version) + sizeof(created)]; my_fseek(file, 4L, MY_SEEK_CUR, MYF(MY_WME)); // skip the event length if (my_fread(file, (byte*) buf, sizeof(buf), MYF(MY_NABP | MY_WME))) return; binlog_version = uint2korr(buf); memcpy(server_version, buf + 2, sizeof(server_version)); created = uint4korr(buf + 2 + sizeof(server_version)); } Start_log_event(const char* buf); ~Start_log_event() {} Log_event_type get_type_code() { return START_EVENT;} int write_data(FILE* file) { if(my_fwrite(file, (byte*) &binlog_version, sizeof(binlog_version), MYF(MY_NABP | MY_WME)) || my_fwrite(file, (byte*) server_version, sizeof(server_version), MYF(MY_NABP | MY_WME)) || my_fwrite(file, (byte*) &created, sizeof(created), MYF(MY_NABP | MY_WME))) return -1; return 0; } int get_data_size() { return sizeof(binlog_version) + sizeof(server_version) + sizeof(created); } void print(FILE* file, bool short_form = 0); }; Loading @@ -295,7 +334,7 @@ class Intvar_log_event: public Log_event Intvar_log_event(uchar type_arg, ulonglong val_arg) :Log_event(time(NULL)),val(val_arg),type(type_arg) {} Intvar_log_event(FILE* file, time_t when); Intvar_log_event(FILE* file, time_t when, uint32 server_id); Intvar_log_event(const char* buf); ~Intvar_log_event() {} Log_event_type get_type_code() { return INTVAR_EVENT;} Loading @@ -311,7 +350,8 @@ class Stop_log_event: public Log_event public: Stop_log_event() :Log_event(time(NULL)) {} Stop_log_event(FILE* file, time_t when_arg): Log_event(when_arg) Stop_log_event(FILE* file, time_t when_arg, uint32 server_id): Log_event(when_arg,0,0,server_id) { my_fseek(file, 4L, MY_SEEK_CUR, MYF(MY_WME)); // skip the event length } Loading @@ -337,7 +377,7 @@ class Rotate_log_event: public Log_event alloced(0) {} Rotate_log_event(FILE* file, time_t when) ; Rotate_log_event(FILE* file, time_t when, uint32 server_id) ; Rotate_log_event(const char* buf, int max_buf); ~Rotate_log_event() { Loading
sql/mysql_priv.h +1 −0 Original line number Diff line number Diff line Loading @@ -437,6 +437,7 @@ void sql_perror(const char *message); void sql_print_error(const char *format,...) __attribute__ ((format (printf, 1, 2))); extern uint32 server_id; extern char mysql_data_home[2],server_version[50],max_sort_char, mysql_real_data_home[]; extern my_string mysql_unix_port,mysql_tmpdir; Loading
sql/mysqlbinlog.cc +2 −0 Original line number Diff line number Diff line Loading @@ -29,6 +29,8 @@ #define CLIENT_CAPABILITIES (CLIENT_LONG_PASSWORD | CLIENT_LONG_FLAG | CLIENT_LOCAL_FILES) char server_version[50]; uint32 server_id = 0; // needed by net_serv.c ulong bytes_sent = 0L, bytes_received = 0L; Loading
sql/mysqld.cc +39 −14 Original line number Diff line number Diff line Loading @@ -142,6 +142,7 @@ static uint handler_count; static bool opt_console=0; #endif static bool opt_skip_slave_start = 0; // if set, slave is not autostarted static ulong opt_specialflag=SPECIAL_ENGLISH; static my_socket unix_sock= INVALID_SOCKET,ip_sock= INVALID_SOCKET; static ulong back_log,connect_timeout,concurrency; Loading Loading @@ -180,6 +181,7 @@ I_List<i_string> replicate_do_db, replicate_ignore_db; // allow the user to tell us which db to replicate and which to ignore I_List<i_string> binlog_do_db, binlog_ignore_db; uint32 server_id = 0; // server id for replication uint mysql_port; uint test_flags, select_errors=0, dropping_tables=0,ha_open_options=0; uint volatile thread_count=0, thread_running=0, kill_cached_threads=0, Loading Loading @@ -1480,6 +1482,8 @@ int main(int argc, char **argv) open_log(&mysql_update_log, hostname, opt_update_logname, "", LOG_NEW); if (opt_bin_log) { if(server_id) { if (!opt_bin_logname) { Loading @@ -1492,6 +1496,10 @@ int main(int argc, char **argv) open_log(&mysql_bin_log, hostname, opt_bin_logname, "-bin", LOG_BIN); } else sql_print_error("Server id is not set - binary logging disabled"); } if (opt_slow_log) open_log(&mysql_slow_log, hostname, opt_slow_logname, "-slow.log", LOG_NORMAL); Loading Loading @@ -1594,11 +1602,16 @@ int main(int argc, char **argv) // slave thread if(master_host) { if(server_id) { pthread_t hThread; if(pthread_create(&hThread, &connection_attrib, handle_slave, 0)) if(!opt_skip_slave_start && pthread_create(&hThread, &connection_attrib, handle_slave, 0)) sql_print_error("Warning: Can't create thread to handle slave"); } else sql_print_error("Server id is not set, slave thread will not be started"); } printf(ER(ER_READY),my_progname,server_version,""); Loading Loading @@ -2201,6 +2214,7 @@ enum options { OPT_LOG_SLAVE_UPDATES, OPT_BINLOG_DO_DB, OPT_BINLOG_IGNORE_DB, OPT_WANT_CORE, OPT_SKIP_CONCURRENT_INSERT, OPT_MEMLOCK, OPT_MYISAM_RECOVER, OPT_REPLICATE_REWRITE_DB, OPT_SERVER_ID, OPT_SKIP_SLAVE_START }; static struct option long_options[] = { Loading Loading @@ -2265,8 +2279,11 @@ static struct option long_options[] = { {"port", required_argument, 0, 'P'}, {"replicate-do-db", required_argument, 0, (int) OPT_REPLICATE_DO_DB}, {"replicate-ignore-db", required_argument, 0, (int) OPT_REPLICATE_IGNORE_DB}, {"replicate-rewrite-db", required_argument, 0, (int) OPT_REPLICATE_REWRITE_DB}, {"safe-mode", no_argument, 0, (int) OPT_SAFE}, {"socket", required_argument, 0, (int) OPT_SOCKET}, {"server-id", required_argument, 0, (int)OPT_SERVER_ID}, {"set-variable", required_argument, 0, 'O'}, #ifdef HAVE_BERKELEY_DB {"skip-bdb", no_argument, 0, (int) OPT_BDB_SKIP}, Loading @@ -2279,6 +2296,7 @@ static struct option long_options[] = { {"skip-name-resolve", no_argument, 0, (int) OPT_SKIP_RESOLVE}, {"skip-new", no_argument, 0, (int) OPT_SKIP_NEW}, {"skip-show-database", no_argument, 0, (int) OPT_SKIP_SHOW_DB}, {"skip-slave-start", no_argument, 0, (int) OPT_SKIP_SLAVE_START}, {"skip-networking", no_argument, 0, (int) OPT_SKIP_NETWORKING}, {"skip-thread-priority", no_argument, 0, (int) OPT_SKIP_PRIOR}, {"sql-bin-update-same", no_argument, 0, (int) OPT_SQL_BIN_UPDATE_SAME}, Loading Loading @@ -2427,6 +2445,7 @@ struct show_var_st init_vars[]= { {"port", (char*) &mysql_port, SHOW_INT}, {"protocol_version", (char*) &protocol_version, SHOW_INT}, {"record_buffer", (char*) &my_default_record_cache_size,SHOW_LONG}, {"server_id", (char*) &server_id, SHOW_LONG}, {"skip_locking", (char*) &my_disable_locking, SHOW_MY_BOOL}, {"skip_networking", (char*) &opt_disable_networking, SHOW_BOOL}, {"skip_show_database", (char*) &opt_skip_show_db, SHOW_BOOL}, Loading Loading @@ -2850,6 +2869,9 @@ static void get_options(int argc,char **argv) opt_slow_log=1; opt_slow_logname=optarg; break; case (int)OPT_SKIP_SLAVE_START: opt_skip_slave_start = 1; break; case (int) OPT_SKIP_NEW: opt_specialflag|= SPECIAL_NO_NEW_FUNC; default_table_type=DB_TYPE_ISAM; Loading Loading @@ -2971,6 +2993,9 @@ static void get_options(int argc,char **argv) default_table_type= (enum db_type) type; break; } case OPT_SERVER_ID: server_id = atoi(optarg); break; case OPT_DELAY_KEY_WRITE: ha_open_options|=HA_OPEN_DELAY_KEY_WRITE; myisam_delay_key_write=1; Loading