log_event.cc 91.5 KB
Newer Older
unknown's avatar
unknown committed
1
/* Copyright (C) 2000-2004 MySQL AB
unknown's avatar
unknown committed
2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22
   
   This program is free software; you can redistribute it and/or modify
   it under the terms of the GNU General Public License as published by
   the Free Software Foundation; either version 2 of the License, or
   (at your option) any later version.
   
   This program is distributed in the hope that it will be useful,
   but WITHOUT ANY WARRANTY; without even the implied warranty of
   MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
   GNU General Public License for more details.
   
   You should have received a copy of the GNU General Public License
   along with this program; if not, write to the Free Software
   Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA  02111-1307  USA */


#ifndef MYSQL_CLIENT
#ifdef __GNUC__
#pragma implementation				// gcc: Class implementation
#endif
#include  "mysql_priv.h"
23
#include "slave.h"
24
#include <my_dir.h>
unknown's avatar
unknown committed
25 26
#endif /* MYSQL_CLIENT */

unknown's avatar
unknown committed
27
#define log_cs	&my_charset_latin1
28

unknown's avatar
unknown committed
29
/*
30
  pretty_print_str()
unknown's avatar
unknown committed
31
*/
32

33
#ifdef MYSQL_CLIENT
34
static void pretty_print_str(FILE* file, char* str, int len)
unknown's avatar
unknown committed
35
{
36
  char* end = str + len;
unknown's avatar
unknown committed
37
  fputc('\'', file);
38 39
  while (str < end)
  {
unknown's avatar
unknown committed
40
    char c;
41 42 43 44 45 46 47 48 49 50 51 52
    switch ((c=*str++)) {
    case '\n': fprintf(file, "\\n"); break;
    case '\r': fprintf(file, "\\r"); break;
    case '\\': fprintf(file, "\\\\"); break;
    case '\b': fprintf(file, "\\b"); break;
    case '\t': fprintf(file, "\\t"); break;
    case '\'': fprintf(file, "\\'"); break;
    case 0   : fprintf(file, "\\0"); break;
    default:
      fputc(c, file);
      break;
    }
53 54
  }
  fputc('\'', file);
unknown's avatar
unknown committed
55
}
unknown's avatar
unknown committed
56
#endif /* MYSQL_CLIENT */
unknown's avatar
unknown committed
57

unknown's avatar
unknown committed
58

unknown's avatar
unknown committed
59
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
unknown's avatar
unknown committed
60

61 62 63 64 65 66 67 68
static void clear_all_errors(THD *thd, struct st_relay_log_info *rli)
{
  thd->query_error = 0;
  thd->clear_error();
  *rli->last_slave_error = 0;
  rli->last_slave_errno = 0;
}

unknown's avatar
unknown committed
69

unknown's avatar
unknown committed
70
/*
unknown's avatar
unknown committed
71
  Ignore error code specified on command line
unknown's avatar
unknown committed
72
*/
73

74 75
inline int ignored_error_code(int err_code)
{
unknown's avatar
unknown committed
76 77
  return ((err_code == ER_SLAVE_IGNORED_TABLE) ||
          (use_slave_mask && bitmap_is_set(&slave_error_mask, err_code)));
78
}
unknown's avatar
SCRUM  
unknown committed
79
#endif
80

unknown's avatar
unknown committed
81

unknown's avatar
unknown committed
82
/*
83
  pretty_print_str()
unknown's avatar
unknown committed
84
*/
unknown's avatar
unknown committed
85

unknown's avatar
unknown committed
86
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
unknown's avatar
unknown committed
87
static char *pretty_print_str(char *packet, char *str, int len)
unknown's avatar
unknown committed
88
{
unknown's avatar
unknown committed
89 90
  char *end= str + len;
  char *pos= packet;
91
  *pos++= '\'';
92 93 94
  while (str < end)
  {
    char c;
95
    switch ((c=*str++)) {
unknown's avatar
unknown committed
96 97 98 99 100 101 102
    case '\n': *pos++= '\\'; *pos++= 'n'; break;
    case '\r': *pos++= '\\'; *pos++= 'r'; break;
    case '\\': *pos++= '\\'; *pos++= '\\'; break;
    case '\b': *pos++= '\\'; *pos++= 'b'; break;
    case '\t': *pos++= '\\'; *pos++= 't'; break;
    case '\'': *pos++= '\\'; *pos++= '\''; break;
    case 0   : *pos++= '\\'; *pos++= '0'; break;
103
    default:
unknown's avatar
unknown committed
104
      *pos++= c;
105 106
      break;
    }
unknown's avatar
unknown committed
107
  }
108 109
  *pos++= '\'';
  return pos;
unknown's avatar
unknown committed
110
}
unknown's avatar
unknown committed
111
#endif /* !MYSQL_CLIENT */
112

unknown's avatar
unknown committed
113

unknown's avatar
unknown committed
114
/*
115
  slave_load_file_stem()
unknown's avatar
unknown committed
116
*/
unknown's avatar
unknown committed
117

unknown's avatar
SCRUM  
unknown committed
118
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
119 120 121
static inline char* slave_load_file_stem(char*buf, uint file_id,
					 int event_server_id)
{
unknown's avatar
unknown committed
122
  fn_format(buf,"SQL_LOAD-",slave_load_tmpdir, "", MY_UNPACK_FILENAME);
123 124 125 126 127 128 129
  buf = strend(buf);
  buf = int10_to_str(::server_id, buf, 10);
  *buf++ = '-';
  buf = int10_to_str(event_server_id, buf, 10);
  *buf++ = '-';
  return int10_to_str(file_id, buf, 10);
}
unknown's avatar
SCRUM  
unknown committed
130
#endif
131

unknown's avatar
unknown committed
132

unknown's avatar
unknown committed
133
/*
134 135
  Delete all temporary files used for SQL_LOAD.

unknown's avatar
unknown committed
136 137
  SYNOPSIS
    cleanup_load_tmpdir()
unknown's avatar
unknown committed
138
*/
139

unknown's avatar
SCRUM  
unknown committed
140
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
141 142 143 144 145
static void cleanup_load_tmpdir()
{
  MY_DIR *dirp;
  FILEINFO *file;
  uint i;
unknown's avatar
unknown committed
146
  char fname[FN_REFLEN], prefbuf[31], *p;
unknown's avatar
unknown committed
147

148 149 150
  if (!(dirp=my_dir(slave_load_tmpdir,MYF(MY_WME))))
    return;

151 152 153 154 155 156 157 158 159 160 161 162 163
  /* 
     When we are deleting temporary files, we should only remove
     the files associated with the server id of our server.
     We don't use event_server_id here because since we've disabled
     direct binlogging of Create_file/Append_file/Exec_load events
     we cannot meet Start_log event in the middle of events from one 
     LOAD DATA.
  */
  p= strmake(prefbuf,"SQL_LOAD-",9);
  p= int10_to_str(::server_id, p, 10);
  *(p++)= '-';
  *p= 0;

164 165 166
  for (i=0 ; i < (uint)dirp->number_off_files; i++)
  {
    file=dirp->dir_entry+i;
167
    if (is_prefix(file->name, prefbuf))
unknown's avatar
unknown committed
168 169 170 171
    {
      fn_format(fname,file->name,slave_load_tmpdir,"",MY_UNPACK_FILENAME);
      my_delete(fname, MYF(0));
    }
172 173 174 175
  }

  my_dirend(dirp);
}
unknown's avatar
SCRUM  
unknown committed
176
#endif
177 178


unknown's avatar
unknown committed
179
/*
180
  write_str()
unknown's avatar
unknown committed
181
*/
182 183 184 185 186 187 188 189

static bool write_str(IO_CACHE *file, char *str, byte length)
{
  return (my_b_safe_write(file, &length, 1) ||
	  my_b_safe_write(file, (byte*) str, (int) length));
}


unknown's avatar
unknown committed
190
/*
191
  read_str()
unknown's avatar
unknown committed
192
*/
193 194 195 196 197 198 199 200 201 202 203 204

static inline int read_str(char * &buf, char *buf_end, char * &str,
			   uint8 &len)
{
  if (buf + (uint) (uchar) *buf >= buf_end)
    return 1;
  len = (uint8) *buf;
  str= buf+1;
  buf+= (uint) len+1;
  return 0;
}

205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227
/*
  Transforms a string into "" or its expression in 0x... form.
*/
static char *str_to_hex(char *to, char *from, uint len)
{
  char *p= to;
  if (len)
  {
    p= strmov(p, "0x");
    for (uint i= 0; i < len; i++, p+= 2)
    {
      /* val[i] is char. Casting to uchar helps greatly if val[i] < 0 */
      uint tmp= (uint) (uchar) from[i];
      p[0]= _dig_vec_upper[tmp >> 4];
      p[1]= _dig_vec_upper[tmp & 15];
    }
    *p= 0;
  }
  else
    p= strmov(p, "\"\"");
  return p; // pointer to end 0 of 'to'
}

228

unknown's avatar
unknown committed
229
/**************************************************************************
unknown's avatar
unknown committed
230
	Log_event methods
unknown's avatar
unknown committed
231
**************************************************************************/
232

unknown's avatar
unknown committed
233
/*
234
  Log_event::get_type_str()
unknown's avatar
unknown committed
235
*/
236

unknown's avatar
unknown committed
237 238
const char* Log_event::get_type_str()
{
239
  switch(get_type_code()) {
unknown's avatar
unknown committed
240 241 242 243 244 245
  case START_EVENT:  return "Start";
  case STOP_EVENT:   return "Stop";
  case QUERY_EVENT:  return "Query";
  case ROTATE_EVENT: return "Rotate";
  case INTVAR_EVENT: return "Intvar";
  case LOAD_EVENT:   return "Load";
246
  case NEW_LOAD_EVENT:   return "New_load";
unknown's avatar
unknown committed
247
  case SLAVE_EVENT:  return "Slave";
248 249 250 251
  case CREATE_FILE_EVENT: return "Create_file";
  case APPEND_BLOCK_EVENT: return "Append_block";
  case DELETE_FILE_EVENT: return "Delete_file";
  case EXEC_LOAD_EVENT: return "Exec_load";
252
  case RAND_EVENT: return "RAND";
unknown's avatar
unknown committed
253
  case USER_VAR_EVENT: return "User var";
unknown's avatar
unknown committed
254
  default: return "Unknown";				/* impossible */ 
unknown's avatar
unknown committed
255 256 257
  }
}

258

unknown's avatar
unknown committed
259
/*
260
  Log_event::Log_event()
unknown's avatar
unknown committed
261
*/
262

unknown's avatar
unknown committed
263
#ifndef MYSQL_CLIENT
264
Log_event::Log_event(THD* thd_arg, uint16 flags_arg, bool using_trans)
unknown's avatar
unknown committed
265
  :log_pos(0), temp_buf(0), exec_time(0), cached_event_len(0),
266
   flags(flags_arg), thd(thd_arg)
267
{
unknown's avatar
unknown committed
268 269 270 271
  server_id=	thd->server_id;
  when=		thd->start_time;
  cache_stmt=	(using_trans &&
		 (thd->options & (OPTION_NOT_AUTOCOMMIT | OPTION_BEGIN)));
272 273 274
}


unknown's avatar
unknown committed
275 276 277 278 279 280 281
/*
  This minimal constructor is for when you are not even sure that there is a
  valid THD. For example in the server when we are shutting down or flushing
  logs after receiving a SIGHUP (then we must write a Rotate to the binlog but
  we have no THD, so we need this minimal constructor).
*/

282 283 284 285
Log_event::Log_event()
  :temp_buf(0), exec_time(0), cached_event_len(0), flags(0), cache_stmt(0),
   thd(0)
{
unknown's avatar
unknown committed
286 287 288
  server_id=	::server_id;
  when=		time(NULL);
  log_pos=	0;
289
}
unknown's avatar
unknown committed
290
#endif /* !MYSQL_CLIENT */
291 292


unknown's avatar
unknown committed
293
/*
294
  Log_event::Log_event()
295
*/
296

297
Log_event::Log_event(const char* buf, bool old_format)
298
  :temp_buf(0), cached_event_len(0), cache_stmt(0)
299 300 301
{
  when = uint4korr(buf);
  server_id = uint4korr(buf + SERVER_ID_OFFSET);
302 303
  if (old_format)
  {
304
    log_pos=0;
305 306 307 308
    flags=0;
  }
  else
  {
309
    log_pos = uint4korr(buf + LOG_POS_OFFSET);
310 311
    flags = uint2korr(buf + FLAGS_OFFSET);
  }
312 313 314 315 316 317
#ifndef MYSQL_CLIENT
  thd = 0;
#endif  
}

#ifndef MYSQL_CLIENT
unknown's avatar
SCRUM  
unknown committed
318
#ifdef HAVE_REPLICATION
319

unknown's avatar
unknown committed
320
/*
321
  Log_event::exec_event()
unknown's avatar
unknown committed
322
*/
323

324
int Log_event::exec_event(struct st_relay_log_info* rli)
325
{
unknown's avatar
unknown committed
326 327
  DBUG_ENTER("Log_event::exec_event");

328 329 330 331 332 333 334 335 336 337 338 339 340
  /*
    rli is null when (as far as I (Guilhem) know)
    the caller is
    Load_log_event::exec_event *and* that one is called from
    Execute_load_log_event::exec_event. 
    In this case, we don't do anything here ;
    Execute_load_log_event::exec_event will call Log_event::exec_event
    again later with the proper rli.
    Strictly speaking, if we were sure that rli is null
    only in the case discussed above, 'if (rli)' is useless here.
    But as we are not 100% sure, keep it for now.
  */
  if (rli)  
341
  {
342 343 344 345 346 347 348 349 350 351 352 353 354 355
    /*
      If in a transaction, and if the slave supports transactions,
      just inc_event_relay_log_pos(). We only have to check for OPTION_BEGIN
      (not OPTION_NOT_AUTOCOMMIT) as transactions are logged
      with BEGIN/COMMIT, not with SET AUTOCOMMIT= .
      
      CAUTION: opt_using_transactions means
      innodb || bdb ; suppose the master supports InnoDB and BDB, 
      but the slave supports only BDB, problems
      will arise: 
      - suppose an InnoDB table is created on the master,
      - then it will be MyISAM on the slave
      - but as opt_using_transactions is true, the slave will believe he is
      transactional with the MyISAM table. And problems will come when one
unknown's avatar
unknown committed
356 357
      does START SLAVE; STOP SLAVE; START SLAVE; (the slave will resume at
      BEGIN whereas there has not been any rollback). This is the problem of
358 359 360 361 362
      using opt_using_transactions instead of a finer
      "does the slave support _the_transactional_handler_used_on_the_master_".
      
      More generally, we'll have problems when a query mixes a transactional
      handler and MyISAM and STOP SLAVE is issued in the middle of the
unknown's avatar
unknown committed
363 364
      "transaction". START SLAVE will resume at BEGIN while the MyISAM table
      has already been updated.
365 366 367
    */
    if ((thd->options & OPTION_BEGIN) && opt_using_transactions)
      rli->inc_event_relay_log_pos(get_event_len());
368 369
    else
    {
370
      rli->inc_group_relay_log_pos(get_event_len(),log_pos);
371
      flush_relay_log_info(rli);
372 373 374 375 376 377
      /* 
         Note that Rotate_log_event::exec_event() does not call this function,
         so there is no chance that a fake rotate event resets
         last_master_timestamp.
      */
      rli->last_master_timestamp= when;
378
    }
379
  }
unknown's avatar
unknown committed
380
  DBUG_RETURN(0);
381
}
unknown's avatar
unknown committed
382

383

unknown's avatar
unknown committed
384
/*
385
  Log_event::pack_info()
unknown's avatar
unknown committed
386
*/
387

388
void Log_event::pack_info(Protocol *protocol)
unknown's avatar
unknown committed
389
{
390
  protocol->store("", &my_charset_bin);
unknown's avatar
unknown committed
391 392 393
}


unknown's avatar
unknown committed
394
/*
395
  Log_event::net_send()
unknown's avatar
unknown committed
396

397
  Only called by SHOW BINLOG EVENTS
unknown's avatar
unknown committed
398
*/
unknown's avatar
SCRUM  
unknown committed
399

400
int Log_event::net_send(Protocol *protocol, const char* log_name, my_off_t pos)
401
{
unknown's avatar
unknown committed
402 403
  const char *p= strrchr(log_name, FN_LIBCHAR);
  const char *event_type;
404 405 406
  if (p)
    log_name = p + 1;
  
407
  protocol->prepare_for_resend();
408
  protocol->store(log_name, &my_charset_bin);
409
  protocol->store((ulonglong) pos);
410
  event_type = get_type_str();
411
  protocol->store(event_type, strlen(event_type), &my_charset_bin);
412 413 414 415
  protocol->store((uint32) server_id);
  protocol->store((ulonglong) log_pos);
  pack_info(protocol);
  return protocol->write();
416
}
unknown's avatar
SCRUM  
unknown committed
417 418 419
#endif /* HAVE_REPLICATION */


unknown's avatar
unknown committed
420
/*
unknown's avatar
SCRUM  
unknown committed
421
  Log_event::init_show_field_list()
unknown's avatar
unknown committed
422
*/
unknown's avatar
SCRUM  
unknown committed
423 424 425 426 427 428 429 430 431 432 433 434 435 436

void Log_event::init_show_field_list(List<Item>* field_list)
{
  field_list->push_back(new Item_empty_string("Log_name", 20));
  field_list->push_back(new Item_return_int("Pos", 11,
					    MYSQL_TYPE_LONGLONG));
  field_list->push_back(new Item_empty_string("Event_type", 20));
  field_list->push_back(new Item_return_int("Server_id", 10,
					    MYSQL_TYPE_LONG));
  field_list->push_back(new Item_return_int("Orig_log_pos", 11,
					    MYSQL_TYPE_LONGLONG));
  field_list->push_back(new Item_empty_string("Info", 20));
}

unknown's avatar
unknown committed
437
#endif /* !MYSQL_CLIENT */
unknown's avatar
unknown committed
438

unknown's avatar
unknown committed
439
/*
440
  Log_event::write()
unknown's avatar
unknown committed
441
*/
442

443
int Log_event::write(IO_CACHE* file)
unknown's avatar
unknown committed
444
{
445
  return (write_header(file) || write_data(file)) ? -1 : 0;
unknown's avatar
unknown committed
446 447
}

448

unknown's avatar
unknown committed
449
/*
450
  Log_event::write_header()
unknown's avatar
unknown committed
451
*/
452

453
int Log_event::write_header(IO_CACHE* file)
unknown's avatar
unknown committed
454
{
455
  char buf[LOG_EVENT_HEADER_LEN];
unknown's avatar
unknown committed
456
  char* pos = buf;
unknown's avatar
unknown committed
457
  int4store(pos, (ulong) when); // timestamp
unknown's avatar
unknown committed
458 459
  pos += 4;
  *pos++ = get_type_code(); // event type code
460 461
  int4store(pos, server_id);
  pos += 4;
462 463
  long tmp=get_data_size() + LOG_EVENT_HEADER_LEN;
  int4store(pos, tmp);
unknown's avatar
unknown committed
464
  pos += 4;
465
  int4store(pos, log_pos);
466 467 468
  pos += 4;
  int2store(pos, flags);
  pos += 2;
469
  return (my_b_safe_write(file, (byte*) buf, (uint) (pos - buf)));
unknown's avatar
unknown committed
470 471 472
}


unknown's avatar
unknown committed
473
/*
474
  Log_event::read_log_event()
unknown's avatar
unknown committed
475
*/
476 477

#ifndef MYSQL_CLIENT
478
int Log_event::read_log_event(IO_CACHE* file, String* packet,
479
			      pthread_mutex_t* log_lock)
unknown's avatar
unknown committed
480 481
{
  ulong data_len;
482
  int result=0;
unknown's avatar
unknown committed
483
  char buf[LOG_EVENT_HEADER_LEN];
484
  DBUG_ENTER("read_log_event");
485

486
  if (log_lock)
487
    pthread_mutex_lock(log_lock);
488 489
  if (my_b_read(file, (byte*) buf, sizeof(buf)))
  {
490 491 492 493 494
    /*
      If the read hits eof, we must report it as eof so the caller
      will know it can go into cond_wait to be woken up on the next
      update to the log.
    */
495
    DBUG_PRINT("error",("file->error: %d", file->error));
496 497 498
    if (!file->error)
      result= LOG_READ_EOF;
    else
499
      result= (file->error > 0 ? LOG_READ_TRUNC : LOG_READ_IO);
500
    goto end;
501
  }
502
  data_len= uint4korr(buf + EVENT_LEN_OFFSET);
unknown's avatar
unknown committed
503 504
  if (data_len < LOG_EVENT_HEADER_LEN ||
      data_len > current_thd->variables.max_allowed_packet)
505
  {
506
    DBUG_PRINT("error",("data_len: %ld", data_len));
507 508 509
    result= ((data_len < LOG_EVENT_HEADER_LEN) ? LOG_READ_BOGUS :
	     LOG_READ_TOO_LARGE);
    goto end;
510
  }
unknown's avatar
unknown committed
511
  packet->append(buf, sizeof(buf));
512
  data_len-= LOG_EVENT_HEADER_LEN;
513 514 515
  if (data_len)
  {
    if (packet->append(file, data_len))
516
    {
517
      /*
518 519
	Here we should never hit EOF in a non-error condition.
	EOF means we are reading the event partially, which should
520 521 522 523
	never happen.
      */
      result= file->error >= 0 ? LOG_READ_TRUNC: LOG_READ_IO;
      /* Implicit goto end; */
524
    }
525
  }
526 527 528 529

end:
  if (log_lock)
    pthread_mutex_unlock(log_lock);
530
  DBUG_RETURN(result);
unknown's avatar
unknown committed
531
}
unknown's avatar
unknown committed
532
#endif /* !MYSQL_CLIENT */
unknown's avatar
unknown committed
533

unknown's avatar
unknown committed
534
#ifndef MYSQL_CLIENT
unknown's avatar
unknown committed
535 536
#define UNLOCK_MUTEX if (log_lock) pthread_mutex_unlock(log_lock);
#define LOCK_MUTEX if (log_lock) pthread_mutex_lock(log_lock);
unknown's avatar
unknown committed
537
#define max_allowed_packet current_thd->variables.max_allowed_packet
538
#else
539
#define UNLOCK_MUTEX
540
#define LOCK_MUTEX
541
#define max_allowed_packet (*mysql_get_parameters()->p_max_allowed_packet)
542 543
#endif

unknown's avatar
unknown committed
544
/*
545 546
  Log_event::read_log_event()

unknown's avatar
unknown committed
547 548 549
  NOTE:
    Allocates memory;  The caller is responsible for clean-up
*/
550

unknown's avatar
unknown committed
551
#ifndef MYSQL_CLIENT
552 553 554
Log_event* Log_event::read_log_event(IO_CACHE* file,
				     pthread_mutex_t* log_lock,
				     bool old_format)
unknown's avatar
unknown committed
555
#else
556
Log_event* Log_event::read_log_event(IO_CACHE* file, bool old_format)
unknown's avatar
unknown committed
557
#endif  
unknown's avatar
unknown committed
558
{
559
  char head[LOG_EVENT_HEADER_LEN];
560 561
  uint header_size= old_format ? OLD_HEADER_LEN : LOG_EVENT_HEADER_LEN;

562
  LOCK_MUTEX;
563
  if (my_b_read(file, (byte *) head, header_size))
564
  {
unknown's avatar
unknown committed
565
    UNLOCK_MUTEX;
566
    return 0;
567
  }
unknown's avatar
unknown committed
568

569
  uint data_len = uint4korr(head + EVENT_LEN_OFFSET);
570 571 572
  char *buf= 0;
  const char *error= 0;
  Log_event *res=  0;
unknown's avatar
unknown committed
573

574
  if (data_len > max_allowed_packet)
unknown's avatar
unknown committed
575
  {
576 577
    error = "Event too big";
    goto err;
unknown's avatar
unknown committed
578 579
  }

580
  if (data_len < header_size)
unknown's avatar
unknown committed
581
  {
582 583
    error = "Event too small";
    goto err;
unknown's avatar
unknown committed
584
  }
585 586 587

  // some events use the extra byte to null-terminate strings
  if (!(buf = my_malloc(data_len+1, MYF(MY_WME))))
588 589 590
  {
    error = "Out of memory";
    goto err;
unknown's avatar
unknown committed
591
  }
592
  buf[data_len] = 0;
593
  memcpy(buf, head, header_size);
594
  if (my_b_read(file, (byte*) buf + header_size, data_len - header_size))
595 596 597 598
  {
    error = "read error";
    goto err;
  }
599
  if ((res = read_log_event(buf, data_len, &error, old_format)))
600
    res->register_temp_buf(buf);
601

602
err:
unknown's avatar
unknown committed
603
  UNLOCK_MUTEX;
604
  if (error)
605
  {
606 607 608
    sql_print_error("\
Error in Log_event::read_log_event(): '%s', data_len: %d, event_type: %d",
		    error,data_len,head[EVENT_TYPE_OFFSET]);
609
    my_free(buf, MYF(MY_ALLOW_ZERO_PTR));
610 611 612 613 614 615 616 617 618
    /*
      The SQL slave thread will check if file->error<0 to know
      if there was an I/O error. Even if there is no "low-level" I/O errors
      with 'file', any of the high-level above errors is worrying
      enough to stop the SQL thread now ; as we are skipping the current event,
      going on with reading and successfully executing other events can
      only corrupt the slave's databases. So stop.
    */
    file->error= -1;
619
  }
620
  return res;
unknown's avatar
unknown committed
621 622
}

623

unknown's avatar
unknown committed
624
/*
625
  Log_event::read_log_event()
unknown's avatar
unknown committed
626
*/
627

628
Log_event* Log_event::read_log_event(const char* buf, int event_len,
629
				     const char **error, bool old_format)
unknown's avatar
unknown committed
630
{
unknown's avatar
unknown committed
631 632
  DBUG_ENTER("Log_event::read_log_event");

633
  if (event_len < EVENT_LEN_OFFSET ||
634 635 636
      (uint) event_len != uint4korr(buf+EVENT_LEN_OFFSET))
  {
    *error="Sanity check failed";		// Needed to free buffer
unknown's avatar
unknown committed
637
    DBUG_RETURN(NULL); // general sanity check - will fail on a partial read
638
  }
639
  
640 641
  Log_event* ev = NULL;
  
642
  switch(buf[EVENT_TYPE_OFFSET]) {
unknown's avatar
unknown committed
643
  case QUERY_EVENT:
644
    ev  = new Query_log_event(buf, event_len, old_format);
645
    break;
unknown's avatar
unknown committed
646
  case LOAD_EVENT:
unknown's avatar
unknown committed
647 648
    ev = new Create_file_log_event(buf, event_len, old_format);
    break;
649
  case NEW_LOAD_EVENT:
650
    ev = new Load_log_event(buf, event_len, old_format);
651
    break;
unknown's avatar
unknown committed
652
  case ROTATE_EVENT:
653
    ev = new Rotate_log_event(buf, event_len, old_format);
654
    break;
unknown's avatar
SCRUM  
unknown committed
655
#ifdef HAVE_REPLICATION
unknown's avatar
unknown committed
656
  case SLAVE_EVENT:
657 658
    ev = new Slave_log_event(buf, event_len);
    break;
unknown's avatar
SCRUM  
unknown committed
659
#endif /* HAVE_REPLICATION */
660
  case CREATE_FILE_EVENT:
unknown's avatar
unknown committed
661
    ev = new Create_file_log_event(buf, event_len, old_format);
662 663 664 665 666 667 668 669 670 671 672
    break;
  case APPEND_BLOCK_EVENT:
    ev = new Append_block_log_event(buf, event_len);
    break;
  case DELETE_FILE_EVENT:
    ev = new Delete_file_log_event(buf, event_len);
    break;
  case EXEC_LOAD_EVENT:
    ev = new Execute_load_log_event(buf, event_len);
    break;
  case START_EVENT:
673
    ev = new Start_log_event(buf, old_format);
674
    break;
unknown's avatar
SCRUM  
unknown committed
675
#ifdef HAVE_REPLICATION
676
  case STOP_EVENT:
677
    ev = new Stop_log_event(buf, old_format);
678
    break;
unknown's avatar
SCRUM  
unknown committed
679
#endif /* HAVE_REPLICATION */
680
  case INTVAR_EVENT:
681
    ev = new Intvar_log_event(buf, old_format);
682
    break;
unknown's avatar
unknown committed
683 684 685
  case RAND_EVENT:
    ev = new Rand_log_event(buf, old_format);
    break;
unknown's avatar
unknown committed
686 687 688
  case USER_VAR_EVENT:
    ev = new User_var_log_event(buf, old_format);
    break;
689 690
  default:
    break;
unknown's avatar
unknown committed
691
  }
692
  if (!ev || !ev->is_valid())
693 694
  {
    delete ev;
695 696 697 698
#ifdef MYSQL_CLIENT
    if (!force_opt)
    {
      *error= "Found invalid event in binary log";
unknown's avatar
unknown committed
699
      DBUG_RETURN(0);
700 701 702 703
    }
    ev= new Unknown_log_event(buf, old_format);
#else
    *error= "Found invalid event in binary log";
unknown's avatar
unknown committed
704
    DBUG_RETURN(0);
705
#endif
706 707
  }
  ev->cached_event_len = event_len;
unknown's avatar
unknown committed
708
  DBUG_RETURN(ev);  
unknown's avatar
unknown committed
709 710
}

711
#ifdef MYSQL_CLIENT
712

unknown's avatar
unknown committed
713
/*
714
  Log_event::print_header()
unknown's avatar
unknown committed
715
*/
716

717 718
void Log_event::print_header(FILE* file)
{
719
  char llbuff[22];
720 721
  fputc('#', file);
  print_timestamp(file);
722
  fprintf(file, " server id %d  log_pos %s ", server_id,
723
	  llstr(log_pos,llbuff)); 
724 725
}

unknown's avatar
unknown committed
726
/*
727
  Log_event::print_timestamp()
unknown's avatar
unknown committed
728
*/
729

730
void Log_event::print_timestamp(FILE* file, time_t* ts)
unknown's avatar
unknown committed
731
{
unknown's avatar
unknown committed
732
  struct tm *res;
733 734
  if (!ts)
    ts = &when;
735 736
#ifdef MYSQL_SERVER				// This is always false
  struct tm tm_tmp;
unknown's avatar
unknown committed
737
  localtime_r(ts,(res= &tm_tmp));
unknown's avatar
unknown committed
738
#else
739
  res=localtime(ts);
unknown's avatar
unknown committed
740
#endif
741 742

  fprintf(file,"%02d%02d%02d %2d:%02d:%02d",
743 744 745 746 747 748
	  res->tm_year % 100,
	  res->tm_mon+1,
	  res->tm_mday,
	  res->tm_hour,
	  res->tm_min,
	  res->tm_sec);
unknown's avatar
unknown committed
749 750
}

unknown's avatar
unknown committed
751
#endif /* MYSQL_CLIENT */
unknown's avatar
unknown committed
752 753


unknown's avatar
unknown committed
754
/*
755
  Log_event::set_log_pos()
unknown's avatar
unknown committed
756
*/
unknown's avatar
unknown committed
757

758 759
#ifndef MYSQL_CLIENT
void Log_event::set_log_pos(MYSQL_LOG* log)
unknown's avatar
unknown committed
760
{
761 762
  if (!log_pos)
    log_pos = my_b_tell(&log->log_file);
unknown's avatar
unknown committed
763
}
unknown's avatar
unknown committed
764
#endif /* !MYSQL_CLIENT */
unknown's avatar
unknown committed
765 766


unknown's avatar
unknown committed
767
/**************************************************************************
unknown's avatar
unknown committed
768
	Query_log_event methods
unknown's avatar
unknown committed
769
**************************************************************************/
770

unknown's avatar
SCRUM  
unknown committed
771
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
772

unknown's avatar
unknown committed
773
/*
774
  Query_log_event::pack_info()
unknown's avatar
unknown committed
775
*/
776

777
void Query_log_event::pack_info(Protocol *protocol)
unknown's avatar
unknown committed
778
{
779 780 781 782
  char *buf, *pos;
  if (!(buf= my_malloc(9 + db_len + q_len, MYF(MY_WME))))
    return;
  pos= buf;    
783
  if (db && db_len)
784
  {
785 786
    pos= strmov(buf, "use `");
    memcpy(pos, db, db_len);
unknown's avatar
unknown committed
787
    pos= strmov(pos+db_len, "`; ");
788
  }
789
  if (query && q_len)
790 791 792 793
  {
    memcpy(pos, query, q_len);
    pos+= q_len;
  }
794
  protocol->store(buf, pos-buf, &my_charset_bin);
795
  my_free(buf, MYF(MY_ALLOW_ZERO_PTR));
796
}
unknown's avatar
SCRUM  
unknown committed
797
#endif
798 799


unknown's avatar
unknown committed
800
/*
801
  Query_log_event::write()
unknown's avatar
unknown committed
802
*/
803 804 805 806

int Query_log_event::write(IO_CACHE* file)
{
  return query ? Log_event::write(file) : -1; 
unknown's avatar
unknown committed
807 808
}

809

unknown's avatar
unknown committed
810
/*
811
  Query_log_event::write_data()
unknown's avatar
unknown committed
812
*/
813 814

int Query_log_event::write_data(IO_CACHE* file)
unknown's avatar
unknown committed
815
{
unknown's avatar
unknown committed
816 817
  char buf[QUERY_HEADER_LEN]; 

818 819 820
  if (!query)
    return -1;
  
unknown's avatar
unknown committed
821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858
  /*
    We want to store the thread id:
    (- as an information for the user when he reads the binlog)
    - if the query uses temporary table: for the slave SQL thread to know to
    which master connection the temp table belongs.
    Now imagine we (write_data()) are called by the slave SQL thread (we are
    logging a query executed by this thread; the slave runs with
    --log-slave-updates). Then this query will be logged with
    thread_id=the_thread_id_of_the_SQL_thread. Imagine that 2 temp tables of
    the same name were created simultaneously on the master (in the master
    binlog you have
    CREATE TEMPORARY TABLE t; (thread 1)
    CREATE TEMPORARY TABLE t; (thread 2)
    ...)
    then in the slave's binlog there will be
    CREATE TEMPORARY TABLE t; (thread_id_of_the_slave_SQL_thread)
    CREATE TEMPORARY TABLE t; (thread_id_of_the_slave_SQL_thread)
    which is bad (same thread id!).

    To avoid this, we log the thread's thread id EXCEPT for the SQL
    slave thread for which we log the original (master's) thread id.
    Now this moves the bug: what happens if the thread id on the
    master was 10 and when the slave replicates the query, a
    connection number 10 is opened by a normal client on the slave,
    and updates a temp table of the same name? We get a problem
    again. To avoid this, in the handling of temp tables (sql_base.cc)
    we use thread_id AND server_id.  TODO when this is merged into
    4.1: in 4.1, slave_proxy_id has been renamed to pseudo_thread_id
    and is a session variable: that's to make mysqlbinlog work with
    temp tables. We probably need to introduce

    SET PSEUDO_SERVER_ID
    for mysqlbinlog in 4.1. mysqlbinlog would print:
    SET PSEUDO_SERVER_ID=
    SET PSEUDO_THREAD_ID=
    for each query using temp tables.
  */
  int4store(buf + Q_THREAD_ID_OFFSET, slave_proxy_id);
859 860 861 862 863 864 865
  int4store(buf + Q_EXEC_TIME_OFFSET, exec_time);
  buf[Q_DB_LEN_OFFSET] = (char) db_len;
  int2store(buf + Q_ERR_CODE_OFFSET, error_code);

  return (my_b_safe_write(file, (byte*) buf, QUERY_HEADER_LEN) ||
	  my_b_safe_write(file, (db) ? (byte*) db : (byte*)"", db_len + 1) ||
	  my_b_safe_write(file, (byte*) query, q_len)) ? -1 : 0;
unknown's avatar
unknown committed
866 867
}

868

unknown's avatar
unknown committed
869
/*
870
  Query_log_event::Query_log_event()
unknown's avatar
unknown committed
871
*/
872

873 874
#ifndef MYSQL_CLIENT
Query_log_event::Query_log_event(THD* thd_arg, const char* query_arg,
875
				 ulong query_length, bool using_trans)
unknown's avatar
unknown committed
876 877
  :Log_event(thd_arg, !thd_arg->tmp_table_used ?
	     0 : LOG_EVENT_THREAD_SPECIFIC_F, using_trans),
unknown's avatar
unknown committed
878
   data_buf(0), query(query_arg),
879
   db(thd_arg->db), q_len((uint32) query_length),
880 881 882
   error_code(thd_arg->killed ?
              ((thd_arg->system_thread & SYSTEM_THREAD_DELAYED_INSERT) ?
               0 : ER_SERVER_SHUTDOWN) : thd_arg->net.last_errno),
unknown's avatar
unknown committed
883 884
   thread_id(thd_arg->thread_id),
   /* save the original thread id; we already know the server id */
885
   slave_proxy_id(thd_arg->variables.pseudo_thread_id)
886 887 888 889 890 891
{
  time_t end_time;
  time(&end_time);
  exec_time = (ulong) (end_time  - thd->start_time);
  db_len = (db) ? (uint32) strlen(db) : 0;
}
unknown's avatar
unknown committed
892
#endif /* MYSQL_CLIENT */
893

unknown's avatar
unknown committed
894

unknown's avatar
unknown committed
895
/*
896
  Query_log_event::Query_log_event()
unknown's avatar
unknown committed
897
*/
898

899
Query_log_event::Query_log_event(const char* buf, int event_len,
unknown's avatar
unknown committed
900 901
				 bool old_format)
  :Log_event(buf, old_format),data_buf(0), query(NULL), db(NULL)
unknown's avatar
unknown committed
902 903
{
  ulong data_len;
904 905 906 907 908 909 910 911 912 913 914 915 916 917
  if (old_format)
  {
    if ((uint)event_len < OLD_HEADER_LEN + QUERY_HEADER_LEN)
      return;				
    data_len = event_len - (QUERY_HEADER_LEN + OLD_HEADER_LEN);
    buf += OLD_HEADER_LEN;
  }
  else
  {
    if ((uint)event_len < QUERY_EVENT_OVERHEAD)
      return;				
    data_len = event_len - QUERY_EVENT_OVERHEAD;
    buf += LOG_EVENT_HEADER_LEN;
  }
unknown's avatar
unknown committed
918

919 920
  exec_time = uint4korr(buf + Q_EXEC_TIME_OFFSET);
  error_code = uint2korr(buf + Q_ERR_CODE_OFFSET);
unknown's avatar
unknown committed
921

922
  if (!(data_buf = (char*) my_malloc(data_len + 1, MYF(MY_WME))))
unknown's avatar
unknown committed
923 924
    return;

925
  memcpy(data_buf, buf + Q_DATA_OFFSET, data_len);
unknown's avatar
unknown committed
926
  slave_proxy_id= thread_id= uint4korr(buf + Q_THREAD_ID_OFFSET);
unknown's avatar
unknown committed
927
  db = data_buf;
928
  db_len = (uint)buf[Q_DB_LEN_OFFSET];
unknown's avatar
unknown committed
929 930 931 932 933
  query=data_buf + db_len + 1;
  q_len = data_len - 1 - db_len;
  *((char*)query+q_len) = 0;
}

934

unknown's avatar
unknown committed
935
/*
936
  Query_log_event::print()
unknown's avatar
unknown committed
937
*/
938

939
#ifdef MYSQL_CLIENT
940
void Query_log_event::print(FILE* file, bool short_form, char* last_db)
unknown's avatar
unknown committed
941
{
942
  char buff[40],*end;				// Enough for SET TIMESTAMP
unknown's avatar
unknown committed
943 944
  if (!short_form)
  {
945
    print_header(file);
946 947
    fprintf(file, "\tQuery\tthread_id=%lu\texec_time=%lu\terror_code=%d\n",
	    (ulong) thread_id, (ulong) exec_time, error_code);
unknown's avatar
unknown committed
948 949
  }

950
  bool different_db= 1;
951

unknown's avatar
unknown committed
952
  if (db && last_db)
953
  {
954
    if (different_db= memcmp(last_db, db, db_len + 1))
955 956
      memcpy(last_db, db, db_len + 1);
  }
957
  
958
  if (db && db[0] && different_db)
unknown's avatar
unknown committed
959
    fprintf(file, "use %s;\n", db);
960 961 962 963
  end=int10_to_str((long) when, strmov(buff,"SET TIMESTAMP="),10);
  *end++=';';
  *end++='\n';
  my_fwrite(file, (byte*) buff, (uint) (end-buff),MYF(MY_NABP | MY_WME));
unknown's avatar
unknown committed
964 965
  if (flags & LOG_EVENT_THREAD_SPECIFIC_F)
    fprintf(file,"SET @@session.pseudo_thread_id=%lu;\n",(ulong)thread_id);
unknown's avatar
unknown committed
966 967 968
  my_fwrite(file, (byte*) query, q_len, MYF(MY_NABP | MY_WME));
  fprintf(file, ";\n");
}
unknown's avatar
unknown committed
969
#endif /* MYSQL_CLIENT */
970

unknown's avatar
unknown committed
971

unknown's avatar
unknown committed
972
/*
973
  Query_log_event::exec_event()
unknown's avatar
unknown committed
974
*/
unknown's avatar
unknown committed
975

unknown's avatar
SCRUM  
unknown committed
976
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
977
int Query_log_event::exec_event(struct st_relay_log_info* rli)
unknown's avatar
unknown committed
978
{
unknown's avatar
unknown committed
979
  int expected_error,actual_error= 0;
unknown's avatar
unknown committed
980
  thd->db= (char*) rewrite_db(db); // thd->db_length is set later if needed
unknown's avatar
unknown committed
981

982
  /*
unknown's avatar
unknown committed
983 984
    InnoDB internally stores the master log position it has processed so far;
    position to store is of the END of the current log event.
985
  */
unknown's avatar
unknown committed
986
#if MYSQL_VERSION_ID < 50000
unknown's avatar
unknown committed
987 988
  rli->future_group_master_log_pos= log_pos + get_event_len() -
    (rli->mi->old_format ? (LOG_EVENT_HEADER_LEN - OLD_HEADER_LEN) : 0);
unknown's avatar
unknown committed
989
#else
unknown's avatar
unknown committed
990
  /* In 5.0 we store the end_log_pos in the relay log so no problem */
unknown's avatar
unknown committed
991 992 993
  rli->future_group_master_log_pos= log_pos;
#endif
  clear_all_errors(thd, rli);
994 995 996 997

  if (db_ok(thd->db, replicate_do_db, replicate_ignore_db))
  {
    thd->set_time((time_t)when);
unknown's avatar
unknown committed
998 999 1000 1001 1002
    /*
      We cannot use db_len from event to fill thd->db_length, because
      rewrite_db() may have changed db.
    */ 
    thd->db_length= thd->db ? strlen(thd->db) : 0;
1003 1004
    thd->query_length= q_len;
    thd->query = (char*)query;
unknown's avatar
unknown committed
1005
    VOID(pthread_mutex_lock(&LOCK_thread_count));
1006 1007
    thd->query_id = query_id++;
    VOID(pthread_mutex_unlock(&LOCK_thread_count));
unknown's avatar
unknown committed
1008
    thd->variables.pseudo_thread_id= thread_id;		// for temp tables
unknown's avatar
unknown committed
1009

unknown's avatar
unknown committed
1010 1011 1012
    mysql_log.write(thd,COM_QUERY,"%s",thd->query);
    DBUG_PRINT("query",("%s",thd->query));
    if (ignored_error_code((expected_error= error_code)) ||
1013 1014
	!check_expected_error(thd,rli,expected_error))
      mysql_parse(thd, thd->query, q_len);
unknown's avatar
unknown committed
1015 1016
    else
    {
unknown's avatar
unknown committed
1017
      /*
unknown's avatar
unknown committed
1018 1019 1020 1021 1022
        The query got a really bad error on the master (thread killed etc),
        which could be inconsistent. Parse it to test the table names: if the
        replicate-*-do|ignore-table rules say "this query must be ignored" then
        we exit gracefully; otherwise we warn about the bad error and tell DBA
        to check/fix it.
unknown's avatar
unknown committed
1023
      */
unknown's avatar
unknown committed
1024 1025 1026
      if (mysql_test_parse_for_slave(thd, thd->query, q_len))
        clear_all_errors(thd, rli);        /* Can ignore query */
      else
1027
      {
unknown's avatar
unknown committed
1028
        slave_print_error(rli,expected_error, 
unknown's avatar
unknown committed
1029
                          "\
unknown's avatar
unknown committed
1030
Query partially completed on the master (error on master: %d) \
unknown's avatar
unknown committed
1031 1032 1033
and was aborted. There is a chance that your master is inconsistent at this \
point. If you are sure that your master is ok, run this query manually on the \
slave and then restart the slave with SET GLOBAL SQL_SLAVE_SKIP_COUNTER=1; \
unknown's avatar
unknown committed
1034
START SLAVE; . Query: '%s'", expected_error, thd->query);
unknown's avatar
unknown committed
1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052
        thd->query_error= 1;
      }
      goto end;
    }

    /*
      If we expected a non-zero error code, and we don't get the same error
      code, and none of them should be ignored.
    */
    DBUG_PRINT("info",("expected_error: %d  last_errno: %d",
		       expected_error, thd->net.last_errno));
    if ((expected_error != (actual_error= thd->net.last_errno)) &&
	expected_error &&
	!ignored_error_code(actual_error) &&
	!ignored_error_code(expected_error))
    {
      slave_print_error(rli, 0,
			"\
unknown's avatar
unknown committed
1053
Query caused different errors on master and slave. \
unknown's avatar
unknown committed
1054
Error on master: '%s' (%d), Error on slave: '%s' (%d). \
unknown's avatar
unknown committed
1055
Default database: '%s'. Query: '%s'",
unknown's avatar
unknown committed
1056 1057 1058 1059
			ER_SAFE(expected_error),
			expected_error,
			actual_error ? thd->net.last_error: "no error",
			actual_error,
unknown's avatar
unknown committed
1060
			print_slave_db_safe(db), query);
unknown's avatar
unknown committed
1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073
      thd->query_error= 1;
    }
    /*
      If we get the same error code as expected, or they should be ignored. 
    */
    else if (expected_error == actual_error ||
	     ignored_error_code(actual_error))
    {
      DBUG_PRINT("info",("error ignored"));
      clear_all_errors(thd, rli);
    }
    /*
      Other cases: mostly we expected no error and get one.
unknown's avatar
unknown committed
1074
    */
unknown's avatar
unknown committed
1075 1076 1077
    else if (thd->query_error || thd->is_fatal_error)
    {
      slave_print_error(rli,actual_error,
unknown's avatar
unknown committed
1078
			"Error '%s' on query. Default database: '%s'. Query: '%s'",
unknown's avatar
unknown committed
1079 1080
			(actual_error ? thd->net.last_error :
			 "unexpected success or fatal error"),
unknown's avatar
unknown committed
1081
			print_slave_db_safe(db), query);
unknown's avatar
unknown committed
1082 1083
      thd->query_error= 1;
    }
unknown's avatar
unknown committed
1084 1085
  } /* End of if (db_ok(... */

unknown's avatar
unknown committed
1086
end:
1087
  VOID(pthread_mutex_lock(&LOCK_thread_count));
unknown's avatar
unknown committed
1088
  thd->db= 0;	                        // prevent db from being freed
1089
  thd->query= 0;			// just to be sure
unknown's avatar
unknown committed
1090
  thd->query_length= thd->db_length =0;
1091
  VOID(pthread_mutex_unlock(&LOCK_thread_count));
unknown's avatar
unknown committed
1092
  close_thread_tables(thd);      
unknown's avatar
unknown committed
1093
  free_root(&thd->mem_root,MYF(MY_KEEP_PREALLOC));
1094 1095 1096 1097 1098 1099 1100 1101 1102
  /*
    If there was an error we stop. Otherwise we increment positions. Note that
    we will not increment group* positions if we are just after a SET
    ONE_SHOT, because SET ONE_SHOT should not be separated from its following
    updating query.
  */
  return (thd->query_error ? thd->query_error : 
          (thd->one_shot_set ? (rli->inc_event_relay_log_pos(get_event_len()),0) :
           Log_event::exec_event(rli))); 
unknown's avatar
unknown committed
1103
}
unknown's avatar
SCRUM  
unknown committed
1104
#endif
unknown's avatar
unknown committed
1105

unknown's avatar
unknown committed
1106

unknown's avatar
unknown committed
1107
/**************************************************************************
unknown's avatar
unknown committed
1108
	Start_log_event methods
unknown's avatar
unknown committed
1109
**************************************************************************/
1110

unknown's avatar
unknown committed
1111
/*
1112
  Start_log_event::pack_info()
unknown's avatar
unknown committed
1113
*/
1114

unknown's avatar
SCRUM  
unknown committed
1115
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
1116
void Start_log_event::pack_info(Protocol *protocol)
unknown's avatar
unknown committed
1117
{
1118 1119 1120 1121
  char buf[12 + ST_SERVER_VER_LEN + 14 + 22], *pos;
  pos= strmov(buf, "Server ver: ");
  pos= strmov(pos, server_version);
  pos= strmov(pos, ", Binlog ver: ");
unknown's avatar
unknown committed
1122 1123
  pos= int10_to_str(binlog_version, pos, 10);
  protocol->store(buf, (uint) (pos-buf), &my_charset_bin);
unknown's avatar
unknown committed
1124
}
unknown's avatar
SCRUM  
unknown committed
1125
#endif
1126 1127


unknown's avatar
unknown committed
1128
/*
1129
  Start_log_event::print()
unknown's avatar
unknown committed
1130
*/
unknown's avatar
unknown committed
1131 1132

#ifdef MYSQL_CLIENT
1133
void Start_log_event::print(FILE* file, bool short_form, char* last_db)
unknown's avatar
unknown committed
1134
{
1135 1136 1137 1138 1139 1140
  if (short_form)
    return;

  print_header(file);
  fprintf(file, "\tStart: binlog v %d, server v %s created ", binlog_version,
	  server_version);
unknown's avatar
unknown committed
1141 1142 1143
  print_timestamp(file);
  if (created)
    fprintf(file," at startup");
1144
  fputc('\n', file);
unknown's avatar
unknown committed
1145 1146
  fflush(file);
}
unknown's avatar
unknown committed
1147
#endif /* MYSQL_CLIENT */
unknown's avatar
unknown committed
1148

unknown's avatar
unknown committed
1149
/*
1150
  Start_log_event::Start_log_event()
unknown's avatar
unknown committed
1151
*/
1152 1153 1154 1155

Start_log_event::Start_log_event(const char* buf,
				 bool old_format)
  :Log_event(buf, old_format)
1156
{
1157 1158 1159 1160 1161
  buf += (old_format) ? OLD_HEADER_LEN : LOG_EVENT_HEADER_LEN;
  binlog_version = uint2korr(buf+ST_BINLOG_VER_OFFSET);
  memcpy(server_version, buf+ST_SERVER_VER_OFFSET,
	 ST_SERVER_VER_LEN);
  created = uint4korr(buf+ST_CREATED_OFFSET);
unknown's avatar
unknown committed
1162 1163
}

1164

unknown's avatar
unknown committed
1165
/*
1166
  Start_log_event::write_data()
unknown's avatar
unknown committed
1167
*/
1168

1169
int Start_log_event::write_data(IO_CACHE* file)
1170
{
1171 1172 1173 1174 1175
  char buff[START_HEADER_LEN];
  int2store(buff + ST_BINLOG_VER_OFFSET,binlog_version);
  memcpy(buff + ST_SERVER_VER_OFFSET,server_version,ST_SERVER_VER_LEN);
  int4store(buff + ST_CREATED_OFFSET,created);
  return (my_b_safe_write(file, (byte*) buff, sizeof(buff)) ? -1 : 0);
1176
}
1177

unknown's avatar
unknown committed
1178
/*
1179 1180 1181 1182 1183
  Start_log_event::exec_event()

  The master started

  IMPLEMENTATION
unknown's avatar
unknown committed
1184 1185 1186 1187
    - To handle the case where the master died without having time to write
      DROP TEMPORARY TABLE, DO RELEASE_LOCK (prepared statements' deletion is
      TODO), we clean up all temporary tables that we got, if we are sure we
      can (see below).
1188 1189

  TODO
1190 1191 1192 1193 1194
    - Remove all active user locks.
      Guilhem 2003-06: this is true but not urgent: the worst it can cause is
      the use of a bit of memory for a user lock which will not be used
      anymore. If the user lock is later used, the old one will be released. In
      other words, no deadlock problem.
unknown's avatar
unknown committed
1195 1196
*/

unknown's avatar
SCRUM  
unknown committed
1197
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
1198 1199
int Start_log_event::exec_event(struct st_relay_log_info* rli)
{
1200
  DBUG_ENTER("Start_log_event::exec_event");
unknown's avatar
unknown committed
1201

unknown's avatar
unknown committed
1202 1203 1204 1205 1206
  /*
    If the I/O thread has not started, mi->old_format is BINLOG_FORMAT_CURRENT
    (that's what the MASTER_INFO constructor does), so the test below is not
    perfect at all.
  */
unknown's avatar
unknown committed
1207 1208 1209 1210 1211 1212 1213 1214
  switch (rli->mi->old_format) {
  case BINLOG_FORMAT_CURRENT:
    /* 
       This is 4.x, so a Start_log_event is only at master startup,
       so we are sure the master has restarted and cleared his temp tables.
    */
    close_temporary_tables(thd);
    cleanup_load_tmpdir();
unknown's avatar
unknown committed
1215 1216 1217 1218 1219 1220 1221 1222 1223 1224
    /*
      As a transaction NEVER spans on 2 or more binlogs:
      if we have an active transaction at this point, the master died while
      writing the transaction to the binary log, i.e. while flushing the binlog
      cache to the binlog. As the write was started, the transaction had been
      committed on the master, so we lack of information to replay this
      transaction on the slave; all we can do is stop with error.
    */
    if (thd->options & OPTION_BEGIN)
    {
unknown's avatar
unknown committed
1225 1226
      slave_print_error(rli, 0, "\
Rolling back unfinished transaction (no COMMIT or ROLLBACK) from relay log. \
unknown's avatar
unknown committed
1227
A probable cause is that the master died while writing the transaction to its \
unknown's avatar
unknown committed
1228
binary log.");
unknown's avatar
unknown committed
1229 1230
      return(1);
    }
unknown's avatar
unknown committed
1231 1232
    break;

unknown's avatar
unknown committed
1233
    /* 
unknown's avatar
unknown committed
1234 1235
       Now the older formats; in that case load_tmpdir is cleaned up by the I/O
       thread.
unknown's avatar
unknown committed
1236
    */
unknown's avatar
unknown committed
1237
  case BINLOG_FORMAT_323_LESS_57:
unknown's avatar
unknown committed
1238
    /*
unknown's avatar
unknown committed
1239 1240 1241
      Cannot distinguish a Start_log_event generated at master startup and
      one generated by master FLUSH LOGS, so cannot be sure temp tables
      have to be dropped. So do nothing.
unknown's avatar
unknown committed
1242
    */
unknown's avatar
unknown committed
1243 1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254
    break;
  case BINLOG_FORMAT_323_GEQ_57:
    /*
      Can distinguish, based on the value of 'created',
      which was generated at master startup.
    */
    if (created)
      close_temporary_tables(thd);
    break;
  default:
    /* this case is impossible */
    return 1;
unknown's avatar
unknown committed
1255
  }
unknown's avatar
unknown committed
1256

1257
  DBUG_RETURN(Log_event::exec_event(rli));
1258
}
unknown's avatar
unknown committed
1259
#endif /* defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT) */
1260

unknown's avatar
unknown committed
1261
/**************************************************************************
unknown's avatar
unknown committed
1262
	Load_log_event methods
unknown's avatar
unknown committed
1263
**************************************************************************/
1264

unknown's avatar
unknown committed
1265
/*
1266
  Load_log_event::pack_info()
unknown's avatar
unknown committed
1267
*/
1268

unknown's avatar
SCRUM  
unknown committed
1269
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
1270
void Load_log_event::pack_info(Protocol *protocol)
1271
{
1272 1273 1274 1275 1276 1277
  char *buf, *pos;
  uint buf_len;

  buf_len= 
    5 + db_len + 3 +                        // "use DB; "
    18 + fname_len + 2 +                    // "LOAD DATA INFILE 'file''"
unknown's avatar
unknown committed
1278
    7 +					    // LOCAL
1279
    9 +                                     // " REPLACE or IGNORE "
1280
    13 + table_name_len*2 +                 // "INTO TABLE `table`"
1281 1282 1283 1284 1285 1286 1287 1288
    21 + sql_ex.field_term_len*4 + 2 +      // " FIELDS TERMINATED BY 'str'"
    23 + sql_ex.enclosed_len*4 + 2 +        // " OPTIONALLY ENCLOSED BY 'str'"
    12 + sql_ex.escaped_len*4 + 2 +         // " ESCAPED BY 'str'"
    21 + sql_ex.line_term_len*4 + 2 +       // " FIELDS TERMINATED BY 'str'"
    19 + sql_ex.line_start_len*4 + 2 +      // " LINES STARTING BY 'str'" 
    15 + 22 +                               // " IGNORE xxx  LINES" 
    3 + (num_fields-1)*2 + field_block_len; // " (field1, field2, ...)"

unknown's avatar
unknown committed
1289
  if (!(buf= my_malloc(buf_len, MYF(MY_WME))))
1290 1291
    return;
  pos= buf;
1292
  if (db && db_len)
1293
  {
1294 1295
    pos= strmov(pos, "use `");
    memcpy(pos, db, db_len);
unknown's avatar
unknown committed
1296
    pos= strmov(pos+db_len, "`; ");
1297
  }
1298

unknown's avatar
unknown committed
1299 1300 1301 1302
  pos= strmov(pos, "LOAD DATA ");
  if (check_fname_outside_temp_buf())
    pos= strmov(pos, "LOCAL ");
  pos= strmov(pos, "INFILE '");
1303
  memcpy(pos, fname, fname_len);
unknown's avatar
unknown committed
1304
  pos= strmov(pos+fname_len, "' ");
1305

unknown's avatar
unknown committed
1306
  if (sql_ex.opt_flags & REPLACE_FLAG)
1307
    pos= strmov(pos, " REPLACE ");
unknown's avatar
unknown committed
1308
  else if (sql_ex.opt_flags & IGNORE_FLAG)
1309 1310
    pos= strmov(pos, " IGNORE ");

1311
  pos= strmov(pos ,"INTO TABLE `");
1312 1313 1314
  memcpy(pos, table_name, table_name_len);
  pos+= table_name_len;

unknown's avatar
unknown committed
1315
  /* We have to create all optinal fields as the default is not empty */
1316
  pos= strmov(pos, "` FIELDS TERMINATED BY ");
unknown's avatar
unknown committed
1317 1318 1319 1320 1321
  pos= pretty_print_str(pos, sql_ex.field_term, sql_ex.field_term_len);
  if (sql_ex.opt_flags & OPT_ENCLOSED_FLAG)
    pos= strmov(pos, " OPTIONALLY ");
  pos= strmov(pos, " ENCLOSED BY ");
  pos= pretty_print_str(pos, sql_ex.enclosed, sql_ex.enclosed_len);
1322

unknown's avatar
unknown committed
1323 1324
  pos= strmov(pos, " ESCAPED BY ");
  pos= pretty_print_str(pos, sql_ex.escaped, sql_ex.escaped_len);
1325

unknown's avatar
unknown committed
1326 1327
  pos= strmov(pos, " LINES TERMINATED BY ");
  pos= pretty_print_str(pos, sql_ex.line_term, sql_ex.line_term_len);
1328 1329
  if (sql_ex.line_start_len)
  {
1330
    pos= strmov(pos, " STARTING BY ");
1331
    pos= pretty_print_str(pos, sql_ex.line_start, sql_ex.line_start_len);
1332
  }
1333

unknown's avatar
unknown committed
1334
  if ((long) skip_lines > 0)
1335 1336
  {
    pos= strmov(pos, " IGNORE ");
unknown's avatar
unknown committed
1337
    pos= longlong10_to_str((longlong) skip_lines, pos, 10);
1338 1339
    pos= strmov(pos," LINES ");    
  }
1340 1341 1342 1343

  if (num_fields)
  {
    uint i;
unknown's avatar
unknown committed
1344
    const char *field= fields;
1345
    pos= strmov(pos, " (");
1346 1347 1348
    for (i = 0; i < num_fields; i++)
    {
      if (i)
unknown's avatar
unknown committed
1349 1350 1351 1352
      {
        *pos++= ' ';
        *pos++= ',';
      }
1353
      memcpy(pos, field, field_lens[i]);
unknown's avatar
unknown committed
1354 1355
      pos+=   field_lens[i];
      field+= field_lens[i]  + 1;
1356
    }
1357
    *pos++= ')';
1358
  }
1359

1360
  protocol->store(buf, pos-buf, &my_charset_bin);
unknown's avatar
unknown committed
1361
  my_free(buf, MYF(0));
1362
}
unknown's avatar
unknown committed
1363
#endif /* defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT) */
1364

1365

unknown's avatar
unknown committed
1366
/*
1367
  Load_log_event::write_data_header()
unknown's avatar
unknown committed
1368
*/
1369 1370

int Load_log_event::write_data_header(IO_CACHE* file)
1371
{
1372
  char buf[LOAD_HEADER_LEN];
unknown's avatar
unknown committed
1373
  int4store(buf + L_THREAD_ID_OFFSET, slave_proxy_id);
1374 1375 1376 1377 1378 1379
  int4store(buf + L_EXEC_TIME_OFFSET, exec_time);
  int4store(buf + L_SKIP_LINES_OFFSET, skip_lines);
  buf[L_TBL_LEN_OFFSET] = (char)table_name_len;
  buf[L_DB_LEN_OFFSET] = (char)db_len;
  int4store(buf + L_NUM_FIELDS_OFFSET, num_fields);
  return my_b_safe_write(file, (byte*)buf, LOAD_HEADER_LEN);
1380
}
1381

1382

unknown's avatar
unknown committed
1383
/*
1384
  Load_log_event::write_data_body()
unknown's avatar
unknown committed
1385
*/
1386 1387

int Load_log_event::write_data_body(IO_CACHE* file)
1388
{
1389 1390 1391
  if (sql_ex.write_data(file))
    return 1;
  if (num_fields && fields && field_lens)
1392
  {
1393 1394 1395
    if (my_b_safe_write(file, (byte*)field_lens, num_fields) ||
	my_b_safe_write(file, (byte*)fields, field_block_len))
      return 1;
1396
  }
1397 1398 1399
  return (my_b_safe_write(file, (byte*)table_name, table_name_len + 1) ||
	  my_b_safe_write(file, (byte*)db, db_len + 1) ||
	  my_b_safe_write(file, (byte*)fname, fname_len));
1400 1401
}

1402

unknown's avatar
unknown committed
1403
/*
1404
  Load_log_event::Load_log_event()
unknown's avatar
unknown committed
1405
*/
1406

1407
#ifndef MYSQL_CLIENT
unknown's avatar
unknown committed
1408 1409 1410
Load_log_event::Load_log_event(THD *thd_arg, sql_exchange *ex,
			       const char *db_arg, const char *table_name_arg,
			       List<Item> &fields_arg,
unknown's avatar
unknown committed
1411 1412 1413
			       enum enum_duplicates handle_dup,
			       bool using_trans)
  :Log_event(thd_arg, 0, using_trans), thread_id(thd_arg->thread_id),
1414
   slave_proxy_id(thd_arg->variables.pseudo_thread_id),
unknown's avatar
unknown committed
1415 1416
   num_fields(0),fields(0),
   field_lens(0),field_block_len(0),
unknown's avatar
unknown committed
1417
   table_name(table_name_arg ? table_name_arg : ""),
1418
   db(db_arg), fname(ex->file_name), local_fname(FALSE)
unknown's avatar
unknown committed
1419 1420 1421
{
  time_t end_time;
  time(&end_time);
1422
  exec_time = (ulong) (end_time  - thd_arg->start_time);
1423 1424 1425
  /* db can never be a zero pointer in 4.0 */
  db_len = (uint32) strlen(db);
  table_name_len = (uint32) strlen(table_name);
unknown's avatar
unknown committed
1426 1427 1428 1429 1430 1431 1432 1433 1434 1435 1436 1437 1438
  fname_len = (fname) ? (uint) strlen(fname) : 0;
  sql_ex.field_term = (char*) ex->field_term->ptr();
  sql_ex.field_term_len = (uint8) ex->field_term->length();
  sql_ex.enclosed = (char*) ex->enclosed->ptr();
  sql_ex.enclosed_len = (uint8) ex->enclosed->length();
  sql_ex.line_term = (char*) ex->line_term->ptr();
  sql_ex.line_term_len = (uint8) ex->line_term->length();
  sql_ex.line_start = (char*) ex->line_start->ptr();
  sql_ex.line_start_len = (uint8) ex->line_start->length();
  sql_ex.escaped = (char*) ex->escaped->ptr();
  sql_ex.escaped_len = (uint8) ex->escaped->length();
  sql_ex.opt_flags = 0;
  sql_ex.cached_new_format = -1;
1439
    
unknown's avatar
unknown committed
1440
  if (ex->dumpfile)
unknown's avatar
unknown committed
1441
    sql_ex.opt_flags|= DUMPFILE_FLAG;
unknown's avatar
unknown committed
1442
  if (ex->opt_enclosed)
unknown's avatar
unknown committed
1443
    sql_ex.opt_flags|= OPT_ENCLOSED_FLAG;
1444

unknown's avatar
unknown committed
1445
  sql_ex.empty_flags= 0;
1446

1447
  switch (handle_dup) {
unknown's avatar
unknown committed
1448
  case DUP_IGNORE:
unknown's avatar
unknown committed
1449
    sql_ex.opt_flags|= IGNORE_FLAG;
unknown's avatar
unknown committed
1450 1451
    break;
  case DUP_REPLACE:
unknown's avatar
unknown committed
1452
    sql_ex.opt_flags|= REPLACE_FLAG;
unknown's avatar
unknown committed
1453 1454 1455 1456
    break;
  case DUP_UPDATE:				// Impossible here
  case DUP_ERROR:
    break;	
unknown's avatar
unknown committed
1457
  }
1458

unknown's avatar
unknown committed
1459 1460 1461 1462 1463 1464 1465 1466 1467 1468
  if (!ex->field_term->length())
    sql_ex.empty_flags |= FIELD_TERM_EMPTY;
  if (!ex->enclosed->length())
    sql_ex.empty_flags |= ENCLOSED_EMPTY;
  if (!ex->line_term->length())
    sql_ex.empty_flags |= LINE_TERM_EMPTY;
  if (!ex->line_start->length())
    sql_ex.empty_flags |= LINE_START_EMPTY;
  if (!ex->escaped->length())
    sql_ex.empty_flags |= ESCAPED_EMPTY;
1469
    
unknown's avatar
unknown committed
1470
  skip_lines = ex->skip_lines;
1471

unknown's avatar
unknown committed
1472 1473 1474 1475 1476 1477 1478 1479 1480 1481 1482
  List_iterator<Item> li(fields_arg);
  field_lens_buf.length(0);
  fields_buf.length(0);
  Item* item;
  while ((item = li++))
  {
    num_fields++;
    uchar len = (uchar) strlen(item->name);
    field_block_len += len + 1;
    fields_buf.append(item->name, len + 1);
    field_lens_buf.append((char*)&len, 1);
1483 1484
  }

unknown's avatar
unknown committed
1485 1486 1487
  field_lens = (const uchar*)field_lens_buf.ptr();
  fields = fields_buf.ptr();
}
unknown's avatar
unknown committed
1488
#endif /* !MYSQL_CLIENT */
1489

1490
/*
1491
  Load_log_event::Load_log_event()
1492

unknown's avatar
unknown committed
1493 1494 1495
  NOTE
    The caller must do buf[event_len] = 0 before he starts using the
    constructed event.
1496 1497
*/

unknown's avatar
unknown committed
1498
Load_log_event::Load_log_event(const char *buf, int event_len,
1499
			       bool old_format)
unknown's avatar
unknown committed
1500 1501 1502
  :Log_event(buf, old_format), num_fields(0), fields(0),
   field_lens(0), field_block_len(0),
   table_name(0), db(0), fname(0), local_fname(FALSE)
unknown's avatar
unknown committed
1503
{
unknown's avatar
unknown committed
1504 1505 1506 1507
  DBUG_ENTER("Load_log_event");
  if (event_len) // derived class, will call copy_log_event() itself
    copy_log_event(buf, event_len, old_format);
  DBUG_VOID_RETURN;
1508 1509
}

1510

unknown's avatar
unknown committed
1511
/*
1512
  Load_log_event::copy_log_event()
unknown's avatar
unknown committed
1513
*/
1514

1515 1516
int Load_log_event::copy_log_event(const char *buf, ulong event_len,
				   bool old_format)
1517
{
1518
  uint data_len;
1519
  char* buf_end = (char*)buf + event_len;
1520
  uint header_len= old_format ? OLD_HEADER_LEN : LOG_EVENT_HEADER_LEN;
1521
  const char* data_head = buf + header_len;
unknown's avatar
unknown committed
1522 1523
  DBUG_ENTER("Load_log_event::copy_log_event");

unknown's avatar
unknown committed
1524
  slave_proxy_id= thread_id= uint4korr(data_head + L_THREAD_ID_OFFSET);
1525 1526 1527 1528 1529
  exec_time = uint4korr(data_head + L_EXEC_TIME_OFFSET);
  skip_lines = uint4korr(data_head + L_SKIP_LINES_OFFSET);
  table_name_len = (uint)data_head[L_TBL_LEN_OFFSET];
  db_len = (uint)data_head[L_DB_LEN_OFFSET];
  num_fields = uint4korr(data_head + L_NUM_FIELDS_OFFSET);
unknown's avatar
unknown committed
1530
	  
1531
  int body_offset = ((buf[EVENT_TYPE_OFFSET] == LOAD_EVENT) ?
1532
		     LOAD_HEADER_LEN + header_len :
1533
		     get_data_body_offset());
unknown's avatar
unknown committed
1534
  
unknown's avatar
unknown committed
1535
  if ((int) event_len < body_offset)
unknown's avatar
unknown committed
1536
    DBUG_RETURN(1);
1537 1538 1539 1540
  /*
    Sql_ex.init() on success returns the pointer to the first byte after
    the sql_ex structure, which is the start of field lengths array.
  */
1541
  if (!(field_lens=(uchar*)sql_ex.init((char*)buf + body_offset,
unknown's avatar
unknown committed
1542 1543 1544 1545
				       buf_end,
				       buf[EVENT_TYPE_OFFSET] != LOAD_EVENT)))
    DBUG_RETURN(1);

1546
  data_len = event_len - body_offset;
1547
  if (num_fields > data_len) // simple sanity check against corruption
unknown's avatar
unknown committed
1548
    DBUG_RETURN(1);
1549
  for (uint i = 0; i < num_fields; i++)
1550
    field_block_len += (uint)field_lens[i] + 1;
1551

unknown's avatar
unknown committed
1552 1553 1554 1555
  fields = (char*)field_lens + num_fields;
  table_name  = fields + field_block_len;
  db = table_name + table_name_len + 1;
  fname = db + db_len + 1;
1556 1557
  fname_len = strlen(fname);
  // null termination is accomplished by the caller doing buf[event_len]=0
unknown's avatar
unknown committed
1558
  DBUG_RETURN(0);
unknown's avatar
unknown committed
1559 1560 1561
}


unknown's avatar
unknown committed
1562
/*
1563
  Load_log_event::print()
unknown's avatar
unknown committed
1564
*/
1565 1566

#ifdef MYSQL_CLIENT
1567
void Load_log_event::print(FILE* file, bool short_form, char* last_db)
unknown's avatar
unknown committed
1568 1569 1570 1571
{
  print(file, short_form, last_db, 0);
}

unknown's avatar
unknown committed
1572 1573 1574

void Load_log_event::print(FILE* file, bool short_form, char* last_db,
			   bool commented)
unknown's avatar
unknown committed
1575
{
unknown's avatar
unknown committed
1576
  DBUG_ENTER("Load_log_event::print");
unknown's avatar
unknown committed
1577 1578
  if (!short_form)
  {
1579
    print_header(file);
1580
    fprintf(file, "\tQuery\tthread_id=%ld\texec_time=%ld\n",
unknown's avatar
unknown committed
1581 1582 1583
	    thread_id, exec_time);
  }

1584
  bool different_db= 1;
unknown's avatar
unknown committed
1585 1586
  if (db && last_db)
  {
1587 1588 1589 1590 1591 1592 1593 1594
    /*
      If the database is different from the one of the previous statement, we
      need to print the "use" command, and we update the last_db.
      But if commented, the "use" is going to be commented so we should not
      update the last_db.
    */
    if ((different_db= memcmp(last_db, db, db_len + 1)) &&
        !commented)
unknown's avatar
unknown committed
1595 1596
      memcpy(last_db, db, db_len + 1);
  }
1597
  
1598
  if (db && db[0] && different_db)
unknown's avatar
unknown committed
1599 1600 1601
    fprintf(file, "%suse %s;\n", 
            commented ? "# " : "",
            db);
unknown's avatar
unknown committed
1602

unknown's avatar
unknown committed
1603 1604
  fprintf(file, "%sLOAD DATA ",
          commented ? "# " : "");
1605 1606
  if (check_fname_outside_temp_buf())
    fprintf(file, "LOCAL ");
1607
  fprintf(file, "INFILE '%-*s' ", fname_len, fname);
unknown's avatar
unknown committed
1608

unknown's avatar
unknown committed
1609
  if (sql_ex.opt_flags & REPLACE_FLAG)
unknown's avatar
unknown committed
1610
    fprintf(file," REPLACE ");
unknown's avatar
unknown committed
1611
  else if (sql_ex.opt_flags & IGNORE_FLAG)
unknown's avatar
unknown committed
1612 1613
    fprintf(file," IGNORE ");
  
1614
  fprintf(file, "INTO TABLE `%s`", table_name);
unknown's avatar
unknown committed
1615 1616
  fprintf(file, " FIELDS TERMINATED BY ");
  pretty_print_str(file, sql_ex.field_term, sql_ex.field_term_len);
unknown's avatar
unknown committed
1617

unknown's avatar
unknown committed
1618 1619 1620 1621
  if (sql_ex.opt_flags & OPT_ENCLOSED_FLAG)
    fprintf(file," OPTIONALLY ");
  fprintf(file, " ENCLOSED BY ");
  pretty_print_str(file, sql_ex.enclosed, sql_ex.enclosed_len);
unknown's avatar
unknown committed
1622
     
unknown's avatar
unknown committed
1623 1624
  fprintf(file, " ESCAPED BY ");
  pretty_print_str(file, sql_ex.escaped, sql_ex.escaped_len);
unknown's avatar
unknown committed
1625
     
unknown's avatar
unknown committed
1626 1627 1628
  fprintf(file," LINES TERMINATED BY ");
  pretty_print_str(file, sql_ex.line_term, sql_ex.line_term_len);

unknown's avatar
unknown committed
1629

1630
  if (sql_ex.line_start)
unknown's avatar
unknown committed
1631
  {
1632
    fprintf(file," STARTING BY ");
1633
    pretty_print_str(file, sql_ex.line_start, sql_ex.line_start_len);
unknown's avatar
unknown committed
1634
  }
1635 1636
  if ((long) skip_lines > 0)
    fprintf(file, " IGNORE %ld LINES", (long) skip_lines);
unknown's avatar
unknown committed
1637

1638 1639 1640 1641
  if (num_fields)
  {
    uint i;
    const char* field = fields;
1642 1643
    fprintf(file, " (");
    for (i = 0; i < num_fields; i++)
unknown's avatar
unknown committed
1644
    {
unknown's avatar
unknown committed
1645
      if (i)
1646 1647
	fputc(',', file);
      fprintf(file, field);
unknown's avatar
unknown committed
1648
	  
1649
      field += field_lens[i]  + 1;
unknown's avatar
unknown committed
1650
    }
1651 1652
    fputc(')', file);
  }
unknown's avatar
unknown committed
1653 1654

  fprintf(file, ";\n");
unknown's avatar
unknown committed
1655
  DBUG_VOID_RETURN;
unknown's avatar
unknown committed
1656
}
unknown's avatar
unknown committed
1657
#endif /* MYSQL_CLIENT */
1658

1659

unknown's avatar
unknown committed
1660
/*
1661
  Load_log_event::set_fields()
unknown's avatar
unknown committed
1662
*/
1663

1664
#ifndef MYSQL_CLIENT
1665
void Load_log_event::set_fields(List<Item> &field_list)
unknown's avatar
unknown committed
1666 1667
{
  uint i;
unknown's avatar
unknown committed
1668
  const char* field = fields;
1669
  for (i= 0; i < num_fields; i++)
unknown's avatar
unknown committed
1670
  {
1671 1672
    field_list.push_back(new Item_field(db, table_name, field));	  
    field+= field_lens[i]  + 1;
unknown's avatar
unknown committed
1673
  }
unknown's avatar
unknown committed
1674
}
unknown's avatar
unknown committed
1675
#endif /* !MYSQL_CLIENT */
unknown's avatar
unknown committed
1676 1677


unknown's avatar
SCRUM  
unknown committed
1678
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
unknown's avatar
unknown committed
1679 1680
/*
  Does the data loading job when executing a LOAD DATA on the slave
1681

unknown's avatar
unknown committed
1682 1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695
  SYNOPSIS
    Load_log_event::exec_event
      net  
      rli                             
      use_rli_only_for_errors	  - if set to 1, rli is provided to 
                                  Load_log_event::exec_event only for this 
				  function to have RPL_LOG_NAME and 
				  rli->last_slave_error, both being used by 
				  error reports. rli's position advancing
				  is skipped (done by the caller which is
				  Execute_load_log_event::exec_event).
				  - if set to 0, rli is provided for full use,
				  i.e. for error reports and position
				  advancing.
1696

unknown's avatar
unknown committed
1697 1698 1699 1700 1701 1702 1703
  DESCRIPTION
    Does the data loading job when executing a LOAD DATA on the slave
 
  RETURN VALUE
    0           Success                                                 
    1    	Failure
*/
1704

unknown's avatar
unknown committed
1705 1706
int Load_log_event::exec_event(NET* net, struct st_relay_log_info* rli, 
			       bool use_rli_only_for_errors)
1707
{
unknown's avatar
unknown committed
1708
  char *load_data_query= 0;
unknown's avatar
unknown committed
1709
  thd->db= (char*) rewrite_db(db); // thd->db_length is set later if needed
1710
  DBUG_ASSERT(thd->query == 0);
unknown's avatar
unknown committed
1711
  thd->query_length= 0;                         // Should not be needed
unknown's avatar
unknown committed
1712
  thd->query_error= 0;
unknown's avatar
unknown committed
1713
  clear_all_errors(thd, rli);
1714 1715 1716 1717
  /*
    Usually mysql_init_query() is called by mysql_parse(), but we need it here
    as the present method does not call mysql_parse().
  */
unknown's avatar
unknown committed
1718
  mysql_init_query(thd, 0, 0);
unknown's avatar
unknown committed
1719
  if (!use_rli_only_for_errors)
1720
  {
unknown's avatar
unknown committed
1721
#if MYSQL_VERSION_ID < 50000
unknown's avatar
unknown committed
1722 1723
    rli->future_group_master_log_pos= log_pos + get_event_len() -
      (rli->mi->old_format ? (LOG_EVENT_HEADER_LEN - OLD_HEADER_LEN) : 0);
unknown's avatar
unknown committed
1724 1725 1726
#else
    rli->future_group_master_log_pos= log_pos;
#endif
1727
  }
1728

unknown's avatar
unknown committed
1729 1730 1731 1732 1733 1734 1735 1736 1737 1738 1739 1740
  /*
    We test replicate_*_db rules. Note that we have already prepared the file
    to load, even if we are going to ignore and delete it now. So it is
    possible that we did a lot of disk writes for nothing. In other words, a
    big LOAD DATA INFILE on the master will still consume a lot of space on
    the slave (space in the relay log + space of temp files: twice the space
    of the file to load...) even if it will finally be ignored.
    TODO: fix this; this can be done by testing rules in
    Create_file_log_event::exec_event() and then discarding Append_block and
    al. Another way is do the filtering in the I/O thread (more efficient: no
    disk writes at all).
  */
1741
  if (db_ok(thd->db, replicate_do_db, replicate_ignore_db))
1742
  {
1743
    thd->set_time((time_t)when);
unknown's avatar
unknown committed
1744
    thd->db_length= thd->db ? strlen(thd->db) : 0;
1745 1746 1747
    VOID(pthread_mutex_lock(&LOCK_thread_count));
    thd->query_id = query_id++;
    VOID(pthread_mutex_unlock(&LOCK_thread_count));
1748 1749 1750 1751 1752 1753 1754
    /*
      Initing thd->row_count is not necessary in theory as this variable has no
      influence in the case of the slave SQL thread (it is used to generate a
      "data truncated" warning but which is absorbed and never gets to the
      error log); still we init it to avoid a Valgrind message.
    */
    mysql_reset_errors(thd);
1755 1756 1757 1758 1759 1760

    TABLE_LIST tables;
    bzero((char*) &tables,sizeof(tables));
    tables.db = thd->db;
    tables.alias = tables.real_name = (char*)table_name;
    tables.lock_type = TL_WRITE;
unknown's avatar
unknown committed
1761
    tables.updating= 1;
unknown's avatar
unknown committed
1762

1763 1764 1765 1766 1767 1768 1769 1770 1771 1772
    // the table will be opened in mysql_load    
    if (table_rules_on && !tables_ok(thd, &tables))
    {
      // TODO: this is a bug - this needs to be moved to the I/O thread
      if (net)
        skip_load_data_infile(net);
    }
    else
    {
      char llbuff[22];
unknown's avatar
unknown committed
1773
      enum enum_duplicates handle_dup;
unknown's avatar
unknown committed
1774 1775 1776 1777 1778 1779 1780 1781 1782 1783 1784 1785 1786
      /*
        Make a simplified LOAD DATA INFILE query, for the information of the
        user in SHOW PROCESSLIST. Note that db is known in the 'db' column.
      */
      if ((load_data_query= (char *) my_alloca(18 + strlen(fname) + 14 +
                                               strlen(tables.real_name) + 8)))
      {
        thd->query_length= (uint)(strxmov(load_data_query,
                                          "LOAD DATA INFILE '", fname,
                                          "' INTO TABLE `", tables.real_name,
                                          "` <...>", NullS) - load_data_query);
        thd->query= load_data_query;
      }
unknown's avatar
unknown committed
1787 1788
      if (sql_ex.opt_flags & REPLACE_FLAG)
	handle_dup= DUP_REPLACE;
unknown's avatar
unknown committed
1789 1790 1791
      else if (sql_ex.opt_flags & IGNORE_FLAG)
        handle_dup= DUP_IGNORE;
      else
unknown's avatar
unknown committed
1792
      {
unknown's avatar
unknown committed
1793
        /*
unknown's avatar
unknown committed
1794
	  When replication is running fine, if it was DUP_ERROR on the
unknown's avatar
unknown committed
1795 1796 1797 1798
          master then we could choose DUP_IGNORE here, because if DUP_ERROR
          suceeded on master, and data is identical on the master and slave,
          then there should be no uniqueness errors on slave, so DUP_IGNORE is
          the same as DUP_ERROR. But in the unlikely case of uniqueness errors
unknown's avatar
unknown committed
1799 1800 1801
          (because the data on the master and slave happen to be different
	  (user error or bug), we want LOAD DATA to print an error message on
	  the slave to discover the problem.
unknown's avatar
unknown committed
1802 1803 1804 1805 1806

          If reading from net (a 3.23 master), mysql_load() will change this
          to DUP_IGNORE.
        */
        handle_dup= DUP_ERROR;
unknown's avatar
unknown committed
1807
      }
unknown's avatar
unknown committed
1808

unknown's avatar
unknown committed
1809
      sql_exchange ex((char*)fname, sql_ex.opt_flags & DUMPFILE_FLAG);
1810 1811 1812 1813 1814
      String field_term(sql_ex.field_term,sql_ex.field_term_len,log_cs);
      String enclosed(sql_ex.enclosed,sql_ex.enclosed_len,log_cs);
      String line_term(sql_ex.line_term,sql_ex.line_term_len,log_cs);
      String line_start(sql_ex.line_start,sql_ex.line_start_len,log_cs);
      String escaped(sql_ex.escaped,sql_ex.escaped_len, log_cs);
unknown's avatar
unknown committed
1815 1816 1817 1818 1819
      ex.field_term= &field_term;
      ex.enclosed= &enclosed;
      ex.line_term= &line_term;
      ex.line_start= &line_start;
      ex.escaped= &escaped;
1820 1821 1822 1823 1824 1825

      ex.opt_enclosed = (sql_ex.opt_flags & OPT_ENCLOSED_FLAG);
      if (sql_ex.empty_flags & FIELD_TERM_EMPTY)
	ex.field_term->length(0);

      ex.skip_lines = skip_lines;
unknown's avatar
unknown committed
1826 1827
      List<Item> field_list;
      set_fields(field_list);
unknown's avatar
unknown committed
1828
      thd->variables.pseudo_thread_id= thread_id;
1829 1830 1831 1832 1833 1834 1835 1836 1837
      if (net)
      {
	// mysql_load will use thd->net to read the file
	thd->net.vio = net->vio;
	/*
	  Make sure the client does not get confused about the packet sequence
	*/
	thd->net.pkt_nr = net->pkt_nr;
      }
unknown's avatar
unknown committed
1838
      if (mysql_load(thd, &ex, &tables, field_list, handle_dup, net != 0,
1839 1840 1841
		     TL_WRITE))
	thd->query_error = 1;
      if (thd->cuted_fields)
unknown's avatar
unknown committed
1842
      {
unknown's avatar
unknown committed
1843 1844 1845 1846 1847 1848 1849 1850
	/* log_pos is the position of the LOAD event in the master log */
	sql_print_error("\
Slave: load data infile on table '%s' at log position %s in log \
'%s' produced %ld warning(s). Default database: '%s'",
                        (char*) table_name,
                        llstr(log_pos,llbuff), RPL_LOG_NAME, 
			(ulong) thd->cuted_fields,
                        print_slave_db_safe(db));
unknown's avatar
unknown committed
1851
      }
1852 1853 1854
      if (net)
        net->pkt_nr= thd->net.pkt_nr;
    }
1855 1856
  }
  else
1857 1858 1859 1860 1861 1862 1863 1864 1865 1866 1867
  {
    /*
      We will just ask the master to send us /dev/null if we do not
      want to load the data.
      TODO: this a bug - needs to be done in I/O thread
    */
    if (net)
      skip_load_data_infile(net);
  }
	    
  thd->net.vio = 0; 
unknown's avatar
unknown committed
1868 1869 1870
  VOID(pthread_mutex_lock(&LOCK_thread_count));
  thd->db= 0;
  thd->query= 0;
unknown's avatar
unknown committed
1871
  thd->query_length= thd->db_length= 0;
unknown's avatar
unknown committed
1872
  VOID(pthread_mutex_unlock(&LOCK_thread_count));
1873
  close_thread_tables(thd);
unknown's avatar
unknown committed
1874 1875
  if (load_data_query)
    my_afree(load_data_query);
1876 1877
  if (thd->query_error)
  {
unknown's avatar
unknown committed
1878 1879 1880 1881 1882 1883 1884 1885 1886 1887
    /* this err/sql_errno code is copy-paste from send_error() */
    const char *err;
    int sql_errno;
    if ((err=thd->net.last_error)[0])
      sql_errno=thd->net.last_errno;
    else
    {
      sql_errno=ER_UNKNOWN_ERROR;
      err=ER(sql_errno);       
    }
unknown's avatar
unknown committed
1888
    slave_print_error(rli,sql_errno,"\
unknown's avatar
unknown committed
1889
Error '%s' running LOAD DATA INFILE on table '%s'. Default database: '%s'",
unknown's avatar
unknown committed
1890
		      err, (char*)table_name, print_slave_db_safe(db));
unknown's avatar
unknown committed
1891
    free_root(&thd->mem_root,MYF(MY_KEEP_PREALLOC));
1892 1893
    return 1;
  }
unknown's avatar
unknown committed
1894
  free_root(&thd->mem_root,MYF(MY_KEEP_PREALLOC));
1895
	    
1896
  if (thd->is_fatal_error)
1897
  {
unknown's avatar
unknown committed
1898 1899 1900
    slave_print_error(rli,ER_UNKNOWN_ERROR, "\
Fatal error running LOAD DATA INFILE on table '%s'. Default database: '%s'",
		      (char*)table_name, print_slave_db_safe(db));
1901 1902 1903
    return 1;
  }

unknown's avatar
unknown committed
1904
  return ( use_rli_only_for_errors ? 0 : Log_event::exec_event(rli) ); 
1905
}
unknown's avatar
SCRUM  
unknown committed
1906
#endif
1907 1908


unknown's avatar
unknown committed
1909
/**************************************************************************
unknown's avatar
unknown committed
1910
  Rotate_log_event methods
unknown's avatar
unknown committed
1911
**************************************************************************/
1912

unknown's avatar
unknown committed
1913
/*
1914
  Rotate_log_event::pack_info()
unknown's avatar
unknown committed
1915
*/
1916

unknown's avatar
SCRUM  
unknown committed
1917
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
1918
void Rotate_log_event::pack_info(Protocol *protocol)
1919
{
unknown's avatar
unknown committed
1920
  char buf1[256], buf[22];
unknown's avatar
unknown committed
1921
  String tmp(buf1, sizeof(buf1), log_cs);
1922
  tmp.length(0);
unknown's avatar
unknown committed
1923 1924 1925 1926
  tmp.append(new_log_ident, ident_len);
  tmp.append(";pos=");
  tmp.append(llstr(pos,buf));
  protocol->store(tmp.ptr(), tmp.length(), &my_charset_bin);
1927
}
unknown's avatar
SCRUM  
unknown committed
1928
#endif
1929

1930

unknown's avatar
unknown committed
1931
/*
1932
  Rotate_log_event::print()
unknown's avatar
unknown committed
1933
*/
1934 1935 1936

#ifdef MYSQL_CLIENT
void Rotate_log_event::print(FILE* file, bool short_form, char* last_db)
1937
{
1938
  char buf[22];
unknown's avatar
unknown committed
1939
  if (short_form)
1940
    return;
1941

1942
  print_header(file);
1943 1944 1945 1946 1947
  fprintf(file, "\tRotate to ");
  if (new_log_ident)
    my_fwrite(file, (byte*) new_log_ident, (uint)ident_len, 
	      MYF(MY_NABP | MY_WME));
  fprintf(file, "  pos: %s", llstr(pos, buf));
1948
  fputc('\n', file);
1949
  fflush(file);
1950
}
unknown's avatar
unknown committed
1951
#endif /* MYSQL_CLIENT */
1952 1953


unknown's avatar
unknown committed
1954
/*
1955
  Rotate_log_event::Rotate_log_event()
unknown's avatar
unknown committed
1956
*/
1957

1958 1959 1960
Rotate_log_event::Rotate_log_event(const char* buf, int event_len,
				   bool old_format)
  :Log_event(buf, old_format),new_log_ident(NULL),alloced(0)
1961
{
1962 1963 1964
  // The caller will ensure that event_len is what we have at EVENT_LEN_OFFSET
  int header_size = (old_format) ? OLD_HEADER_LEN : LOG_EVENT_HEADER_LEN;
  uint ident_offset;
unknown's avatar
unknown committed
1965 1966
  DBUG_ENTER("Rotate_log_event");

1967
  if (event_len < header_size)
unknown's avatar
unknown committed
1968 1969
    DBUG_VOID_RETURN;

1970 1971 1972 1973 1974 1975
  buf += header_size;
  if (old_format)
  {
    ident_len = (uint)(event_len - OLD_HEADER_LEN);
    pos = 4;
    ident_offset = 0;
1976
  }
1977 1978 1979 1980 1981 1982 1983 1984 1985 1986 1987
  else
  {
    ident_len = (uint)(event_len - ROTATE_EVENT_OVERHEAD);
    pos = uint8korr(buf + R_POS_OFFSET);
    ident_offset = ROTATE_HEADER_LEN;
  }
  set_if_smaller(ident_len,FN_REFLEN-1);
  if (!(new_log_ident= my_strdup_with_length((byte*) buf +
					     ident_offset,
					     (uint) ident_len,
					     MYF(MY_WME))))
unknown's avatar
unknown committed
1988
    DBUG_VOID_RETURN;
1989
  alloced = 1;
unknown's avatar
unknown committed
1990
  DBUG_VOID_RETURN;
1991
}
1992 1993


unknown's avatar
unknown committed
1994
/*
1995
  Rotate_log_event::write_data()
unknown's avatar
unknown committed
1996
*/
1997

1998
int Rotate_log_event::write_data(IO_CACHE* file)
1999
{
2000
  char buf[ROTATE_HEADER_LEN];
unknown's avatar
unknown committed
2001
  int8store(buf + R_POS_OFFSET, pos);
2002 2003
  return (my_b_safe_write(file, (byte*)buf, ROTATE_HEADER_LEN) ||
	  my_b_safe_write(file, (byte*)new_log_ident, (uint) ident_len));
2004 2005
}

2006

unknown's avatar
unknown committed
2007
/*
2008 2009 2010
  Rotate_log_event::exec_event()

  Got a rotate log even from the master
2011

2012 2013 2014
  IMPLEMENTATION
    This is mainly used so that we can later figure out the logname and
    position for the master.
2015

2016 2017 2018 2019 2020
    We can't rotate the slave as this will cause infinitive rotations
    in a A -> B -> A setup.

  RETURN VALUES
    0	ok
unknown's avatar
unknown committed
2021
*/
2022

unknown's avatar
SCRUM  
unknown committed
2023
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
2024
int Rotate_log_event::exec_event(struct st_relay_log_info* rli)
2025
{
2026 2027 2028
  DBUG_ENTER("Rotate_log_event::exec_event");

  pthread_mutex_lock(&rli->data_lock);
unknown's avatar
unknown committed
2029
  rli->event_relay_log_pos += get_event_len();
unknown's avatar
unknown committed
2030 2031 2032 2033 2034 2035 2036 2037 2038
  /*
    If we are in a transaction: the only normal case is when the I/O thread was
    copying a big transaction, then it was stopped and restarted: we have this
    in the relay log:
    BEGIN
    ...
    ROTATE (a fake one)
    ...
    COMMIT or ROLLBACK
unknown's avatar
unknown committed
2039 2040
    In that case, we don't want to touch the coordinates which correspond to
    the beginning of the transaction.
unknown's avatar
unknown committed
2041
  */
unknown's avatar
unknown committed
2042
  if (!(thd->options & OPTION_BEGIN))
unknown's avatar
unknown committed
2043
  {
unknown's avatar
unknown committed
2044 2045 2046 2047 2048 2049
    memcpy(rli->group_master_log_name, new_log_ident, ident_len+1);
    rli->notify_group_master_log_name_update();
    rli->group_master_log_pos = pos;
    rli->group_relay_log_pos = rli->event_relay_log_pos;
    DBUG_PRINT("info", ("group_master_log_pos: %lu",
                        (ulong) rli->group_master_log_pos));
unknown's avatar
unknown committed
2050
  }
2051 2052 2053 2054
  pthread_mutex_unlock(&rli->data_lock);
  pthread_cond_broadcast(&rli->data_cond);
  flush_relay_log_info(rli);
  DBUG_RETURN(0);
2055
}
unknown's avatar
SCRUM  
unknown committed
2056
#endif
2057 2058


unknown's avatar
unknown committed
2059
/**************************************************************************
unknown's avatar
unknown committed
2060
	Intvar_log_event methods
unknown's avatar
unknown committed
2061
**************************************************************************/
2062

unknown's avatar
unknown committed
2063
/*
2064
  Intvar_log_event::pack_info()
unknown's avatar
unknown committed
2065
*/
2066

unknown's avatar
SCRUM  
unknown committed
2067
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
2068
void Intvar_log_event::pack_info(Protocol *protocol)
2069
{
unknown's avatar
unknown committed
2070 2071
  char buf[256], *pos;
  pos= strmake(buf, get_var_type_name(), sizeof(buf)-23);
unknown's avatar
unknown committed
2072
  *pos++= '=';
unknown's avatar
unknown committed
2073
  pos= longlong10_to_str(val, pos, -10);
unknown's avatar
unknown committed
2074
  protocol->store(buf, (uint) (pos-buf), &my_charset_bin);
2075
}
unknown's avatar
SCRUM  
unknown committed
2076
#endif
2077

unknown's avatar
unknown committed
2078

unknown's avatar
unknown committed
2079
/*
2080
  Intvar_log_event::Intvar_log_event()
unknown's avatar
unknown committed
2081
*/
2082

2083 2084
Intvar_log_event::Intvar_log_event(const char* buf, bool old_format)
  :Log_event(buf, old_format)
2085
{
2086 2087 2088
  buf += (old_format) ? OLD_HEADER_LEN : LOG_EVENT_HEADER_LEN;
  type = buf[I_TYPE_OFFSET];
  val = uint8korr(buf+I_VAL_OFFSET);
2089 2090
}

2091

unknown's avatar
unknown committed
2092
/*
2093
  Intvar_log_event::get_var_type_name()
unknown's avatar
unknown committed
2094
*/
2095 2096

const char* Intvar_log_event::get_var_type_name()
2097
{
2098 2099 2100 2101 2102
  switch(type) {
  case LAST_INSERT_ID_EVENT: return "LAST_INSERT_ID";
  case INSERT_ID_EVENT: return "INSERT_ID";
  default: /* impossible */ return "UNKNOWN";
  }
2103 2104
}

unknown's avatar
unknown committed
2105

unknown's avatar
unknown committed
2106
/*
2107
  Intvar_log_event::write_data()
unknown's avatar
unknown committed
2108
*/
2109 2110

int Intvar_log_event::write_data(IO_CACHE* file)
2111
{
2112 2113 2114 2115
  char buf[9];
  buf[I_TYPE_OFFSET] = type;
  int8store(buf + I_VAL_OFFSET, val);
  return my_b_safe_write(file, (byte*) buf, sizeof(buf));
2116 2117
}

2118

unknown's avatar
unknown committed
2119
/*
2120
  Intvar_log_event::print()
unknown's avatar
unknown committed
2121
*/
2122 2123 2124

#ifdef MYSQL_CLIENT
void Intvar_log_event::print(FILE* file, bool short_form, char* last_db)
2125
{
2126 2127 2128
  char llbuff[22];
  const char *msg;
  LINT_INIT(msg);
2129

2130 2131 2132 2133 2134
  if (!short_form)
  {
    print_header(file);
    fprintf(file, "\tIntvar\n");
  }
2135

2136 2137 2138 2139 2140 2141 2142 2143 2144 2145 2146
  fprintf(file, "SET ");
  switch (type) {
  case LAST_INSERT_ID_EVENT:
    msg="LAST_INSERT_ID";
    break;
  case INSERT_ID_EVENT:
    msg="INSERT_ID";
    break;
  }
  fprintf(file, "%s=%s;\n", msg, llstr(val,llbuff));
  fflush(file);
2147
}
2148
#endif
2149

2150

unknown's avatar
unknown committed
2151
/*
2152
  Intvar_log_event::exec_event()
unknown's avatar
unknown committed
2153
*/
2154

unknown's avatar
SCRUM  
unknown committed
2155
#if defined(HAVE_REPLICATION)&& !defined(MYSQL_CLIENT)
2156
int Intvar_log_event::exec_event(struct st_relay_log_info* rli)
2157
{
2158 2159 2160 2161 2162 2163 2164 2165 2166
  switch (type) {
  case LAST_INSERT_ID_EVENT:
    thd->last_insert_id_used = 1;
    thd->last_insert_id = val;
    break;
  case INSERT_ID_EVENT:
    thd->next_insert_id = val;
    break;
  }
2167
  rli->inc_event_relay_log_pos(get_event_len());
2168
  return 0;
2169
}
unknown's avatar
SCRUM  
unknown committed
2170
#endif
2171

2172

unknown's avatar
unknown committed
2173
/**************************************************************************
2174
  Rand_log_event methods
unknown's avatar
unknown committed
2175
**************************************************************************/
2176

unknown's avatar
SCRUM  
unknown committed
2177
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
2178
void Rand_log_event::pack_info(Protocol *protocol)
2179
{
unknown's avatar
unknown committed
2180 2181 2182 2183 2184
  char buf1[256], *pos;
  pos= strmov(buf1,"rand_seed1=");
  pos= int10_to_str((long) seed1, pos, 10);
  pos= strmov(pos, ",rand_seed2=");
  pos= int10_to_str((long) seed2, pos, 10);
2185
  protocol->store(buf1, (uint) (pos-buf1), &my_charset_bin);
2186
}
unknown's avatar
SCRUM  
unknown committed
2187
#endif
2188 2189 2190 2191


Rand_log_event::Rand_log_event(const char* buf, bool old_format)
  :Log_event(buf, old_format)
2192
{
2193 2194 2195
  buf += (old_format) ? OLD_HEADER_LEN : LOG_EVENT_HEADER_LEN;
  seed1 = uint8korr(buf+RAND_SEED1_OFFSET);
  seed2 = uint8korr(buf+RAND_SEED2_OFFSET);
2196 2197
}

2198 2199

int Rand_log_event::write_data(IO_CACHE* file)
2200
{
2201 2202 2203 2204 2205
  char buf[16];
  int8store(buf + RAND_SEED1_OFFSET, seed1);
  int8store(buf + RAND_SEED2_OFFSET, seed2);
  return my_b_safe_write(file, (byte*) buf, sizeof(buf));
}
2206

2207 2208 2209 2210

#ifdef MYSQL_CLIENT
void Rand_log_event::print(FILE* file, bool short_form, char* last_db)
{
unknown's avatar
unknown committed
2211
  char llbuff[22],llbuff2[22];
2212
  if (!short_form)
2213
  {
2214 2215
    print_header(file);
    fprintf(file, "\tRand\n");
2216
  }
unknown's avatar
unknown committed
2217
  fprintf(file, "SET @@RAND_SEED1=%s, @@RAND_SEED2=%s;\n",
unknown's avatar
unknown committed
2218
	  llstr(seed1, llbuff),llstr(seed2, llbuff2));
2219
  fflush(file);
2220
}
unknown's avatar
unknown committed
2221
#endif /* MYSQL_CLIENT */
2222

2223

unknown's avatar
SCRUM  
unknown committed
2224
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
2225
int Rand_log_event::exec_event(struct st_relay_log_info* rli)
2226
{
unknown's avatar
unknown committed
2227 2228
  thd->rand.seed1= (ulong) seed1;
  thd->rand.seed2= (ulong) seed2;
2229
  rli->inc_event_relay_log_pos(get_event_len());
2230 2231
  return 0;
}
unknown's avatar
unknown committed
2232
#endif /* !MYSQL_CLIENT */
2233

unknown's avatar
unknown committed
2234

unknown's avatar
unknown committed
2235
/**************************************************************************
2236
  User_var_log_event methods
unknown's avatar
unknown committed
2237
**************************************************************************/
unknown's avatar
unknown committed
2238

unknown's avatar
unknown committed
2239
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
unknown's avatar
unknown committed
2240 2241 2242
void User_var_log_event::pack_info(Protocol* protocol)
{
  char *buf= 0;
2243
  uint val_offset= 4 + name_len;
unknown's avatar
unknown committed
2244 2245 2246 2247 2248 2249 2250 2251 2252 2253 2254 2255 2256 2257 2258
  uint event_len= val_offset;

  if (is_null)
  {
    buf= my_malloc(val_offset + 5, MYF(MY_WME));
    strmov(buf + val_offset, "NULL");
    event_len= val_offset + 4;
  }
  else
  {
    switch (type) {
    case REAL_RESULT:
      double real_val;
      float8get(real_val, val);
      buf= my_malloc(val_offset + FLOATING_POINT_BUFFER, MYF(MY_WME));
2259 2260
      event_len+= my_sprintf(buf + val_offset,
			     (buf + val_offset, "%.14g", real_val));
unknown's avatar
unknown committed
2261 2262 2263 2264 2265 2266
      break;
    case INT_RESULT:
      buf= my_malloc(val_offset + 22, MYF(MY_WME));
      event_len= longlong10_to_str(uint8korr(val), buf + val_offset,-10)-buf;
      break;
    case STRING_RESULT:
2267 2268 2269 2270 2271 2272 2273 2274 2275 2276
      /* 15 is for 'COLLATE' and other chars */
      buf= my_malloc(event_len+val_len*2+1+2*MY_CS_NAME_SIZE+15, MYF(MY_WME));
      CHARSET_INFO *cs;
      if (!(cs= get_charset(charset_number, MYF(0))))
      {
        strmov(buf+val_offset, "???");
        event_len+= 3;
      }
      else
      {
2277 2278 2279
        char *p= strxmov(buf + val_offset, "_", cs->csname, " ", NullS);
        p= str_to_hex(p, val, val_len);
        p= strxmov(p, " COLLATE ", cs->name, NullS);
2280 2281
        event_len= p-buf;
      }
unknown's avatar
unknown committed
2282
      break;
2283
    case ROW_RESULT:
unknown's avatar
unknown committed
2284
    default:
unknown's avatar
unknown committed
2285 2286 2287 2288 2289
      DBUG_ASSERT(1);
      return;
    }
  }
  buf[0]= '@';
2290 2291 2292 2293
  buf[1]= '`';
  buf[2+name_len]= '`';
  buf[3+name_len]= '=';
  memcpy(buf+2, name, name_len);
2294
  protocol->store(buf, event_len, &my_charset_bin);
unknown's avatar
unknown committed
2295 2296
  my_free(buf, MYF(MY_ALLOW_ZERO_PTR));
}
unknown's avatar
unknown committed
2297
#endif /* !MYSQL_CLIENT */
unknown's avatar
unknown committed
2298 2299 2300 2301 2302 2303 2304 2305


User_var_log_event::User_var_log_event(const char* buf, bool old_format)
  :Log_event(buf, old_format)
{
  buf+= (old_format) ? OLD_HEADER_LEN : LOG_EVENT_HEADER_LEN;
  name_len= uint4korr(buf);
  name= (char *) buf + UV_NAME_LEN_SIZE;
2306 2307
  buf+= UV_NAME_LEN_SIZE + name_len;
  is_null= (bool) *buf;
unknown's avatar
unknown committed
2308 2309 2310
  if (is_null)
  {
    type= STRING_RESULT;
2311
    charset_number= my_charset_bin.number;
unknown's avatar
unknown committed
2312 2313 2314 2315 2316
    val_len= 0;
    val= 0;  
  }
  else
  {
2317 2318 2319
    type= (Item_result) buf[UV_VAL_IS_NULL];
    charset_number= uint4korr(buf + UV_VAL_IS_NULL + UV_VAL_TYPE_SIZE);
    val_len= uint4korr(buf + UV_VAL_IS_NULL + UV_VAL_TYPE_SIZE + 
unknown's avatar
unknown committed
2320
		       UV_CHARSET_NUMBER_SIZE);
2321 2322
    val= (char *) (buf + UV_VAL_IS_NULL + UV_VAL_TYPE_SIZE +
		   UV_CHARSET_NUMBER_SIZE + UV_VAL_LEN_SIZE);
unknown's avatar
unknown committed
2323 2324 2325 2326 2327 2328 2329 2330 2331
  }
}


int User_var_log_event::write_data(IO_CACHE* file)
{
  char buf[UV_NAME_LEN_SIZE];
  char buf1[UV_VAL_IS_NULL + UV_VAL_TYPE_SIZE + 
	    UV_CHARSET_NUMBER_SIZE + UV_VAL_LEN_SIZE];
2332 2333 2334
  char buf2[8], *pos= buf2;
  uint buf1_length;

unknown's avatar
unknown committed
2335
  int4store(buf, name_len);
2336 2337 2338 2339 2340 2341 2342
  
  if ((buf1[0]= is_null))
  {
    buf1_length= 1;
    val_len= 0;
  }    
  else
unknown's avatar
unknown committed
2343 2344 2345 2346
  {
    buf1[1]= type;
    int4store(buf1 + 2, charset_number);
    int4store(buf1 + 2 + UV_CHARSET_NUMBER_SIZE, val_len);
2347
    buf1_length= 10;
unknown's avatar
unknown committed
2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358

    switch (type) {
    case REAL_RESULT:
      float8store(buf2, *(double*) val);
      break;
    case INT_RESULT:
      int8store(buf2, *(longlong*) val);
      break;
    case STRING_RESULT:
      pos= val;
      break;
2359
    case ROW_RESULT:
unknown's avatar
unknown committed
2360
    default:
unknown's avatar
unknown committed
2361 2362 2363 2364 2365
      DBUG_ASSERT(1);
      return 0;
    }
  }
  return (my_b_safe_write(file, (byte*) buf, sizeof(buf))   ||
2366 2367 2368
	  my_b_safe_write(file, (byte*) name, name_len)     ||
	  my_b_safe_write(file, (byte*) buf1, buf1_length) ||
	  my_b_safe_write(file, (byte*) pos, val_len));
unknown's avatar
unknown committed
2369 2370
}

2371

unknown's avatar
unknown committed
2372
/*
unknown's avatar
unknown committed
2373
  User_var_log_event::print()
unknown's avatar
unknown committed
2374
*/
unknown's avatar
unknown committed
2375 2376 2377 2378 2379 2380 2381 2382 2383 2384

#ifdef MYSQL_CLIENT
void User_var_log_event::print(FILE* file, bool short_form, char* last_db)
{
  if (!short_form)
  {
    print_header(file);
    fprintf(file, "\tUser_var\n");
  }

2385
  fprintf(file, "SET @`");
unknown's avatar
unknown committed
2386
  my_fwrite(file, (byte*) name, (uint) (name_len), MYF(MY_NABP | MY_WME));
2387
  fprintf(file, "`");
unknown's avatar
unknown committed
2388 2389 2390 2391 2392 2393 2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404 2405 2406

  if (is_null)
  {
    fprintf(file, ":=NULL;\n");
  }
  else
  {
    switch (type) {
    case REAL_RESULT:
      double real_val;
      float8get(real_val, val);
      fprintf(file, ":=%.14g;\n", real_val);
      break;
    case INT_RESULT:
      char int_buf[22];
      longlong10_to_str(uint8korr(val), int_buf, -10);
      fprintf(file, ":=%s;\n", int_buf);
      break;
    case STRING_RESULT:
2407
    {
2408 2409 2410 2411 2412 2413 2414 2415 2416 2417 2418 2419 2420 2421
      /*
        Let's express the string in hex. That's the most robust way. If we
        print it in character form instead, we need to escape it with
        character_set_client which we don't know (we will know it in 5.0, but
        in 4.1 we don't know it easily when we are printing
        User_var_log_event). Explanation why we would need to bother with
        character_set_client (quoting Bar):
        > Note, the parser doesn't switch to another unescaping mode after
        > it has met a character set introducer.
        > For example, if an SJIS client says something like:
        > SET @a= _ucs2 \0a\0b'
        > the string constant is still unescaped according to SJIS, not
        > according to UCS2.
      */
2422 2423 2424 2425
      char *hex_str;
      CHARSET_INFO *cs;

      if (!(hex_str= (char *)my_alloca(2*val_len+1+2))) // 2 hex digits / byte
2426
        break; // no error, as we are 'void'
2427
      str_to_hex(hex_str, val, val_len);
2428 2429 2430 2431 2432 2433 2434 2435 2436 2437 2438 2439 2440
      /*
        For proper behaviour when mysqlbinlog|mysql, we need to explicitely
        specify the variable's collation. It will however cause problems when
        people want to mysqlbinlog|mysql into another server not supporting the
        character set. But there's not much to do about this and it's unlikely.
      */
      if (!(cs= get_charset(charset_number, MYF(0))))
        /*
          Generate an unusable command (=> syntax error) is probably the best
          thing we can do here.
        */
        fprintf(file, ":=???;\n");
      else
2441 2442
        fprintf(file, ":=_%s %s COLLATE %s;\n", cs->csname, hex_str, cs->name);
      my_afree(hex_str);
2443
    }
unknown's avatar
unknown committed
2444
      break;
2445
    case ROW_RESULT:
unknown's avatar
unknown committed
2446
    default:
unknown's avatar
unknown committed
2447
      DBUG_ASSERT(1);
unknown's avatar
unknown committed
2448 2449 2450 2451 2452
      return;
    }
  }
  fflush(file);
}
unknown's avatar
SCRUM  
unknown committed
2453
#endif
2454

unknown's avatar
unknown committed
2455

unknown's avatar
unknown committed
2456
/*
unknown's avatar
unknown committed
2457
  User_var_log_event::exec_event()
unknown's avatar
unknown committed
2458
*/
unknown's avatar
unknown committed
2459

unknown's avatar
unknown committed
2460
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
unknown's avatar
unknown committed
2461 2462 2463
int User_var_log_event::exec_event(struct st_relay_log_info* rli)
{
  Item *it= 0;
2464 2465 2466
  CHARSET_INFO *charset;
  if (!(charset= get_charset(charset_number, MYF(MY_WME))))
    return 1;
unknown's avatar
unknown committed
2467 2468 2469
  LEX_STRING user_var_name;
  user_var_name.str= name;
  user_var_name.length= name_len;
2470 2471
  double real_val;
  longlong int_val;
unknown's avatar
unknown committed
2472 2473 2474 2475 2476 2477 2478 2479 2480 2481 2482

  if (is_null)
  {
    it= new Item_null();
  }
  else
  {
    switch (type) {
    case REAL_RESULT:
      float8get(real_val, val);
      it= new Item_real(real_val);
2483
      val= (char*) &real_val;		// Pointer to value in native format
2484
      val_len= 8;
unknown's avatar
unknown committed
2485 2486
      break;
    case INT_RESULT:
2487 2488 2489
      int_val= (longlong) uint8korr(val);
      it= new Item_int(int_val);
      val= (char*) &int_val;		// Pointer to value in native format
2490
      val_len= 8;
unknown's avatar
unknown committed
2491 2492 2493 2494
      break;
    case STRING_RESULT:
      it= new Item_string(val, val_len, charset);
      break;
2495
    case ROW_RESULT:
unknown's avatar
unknown committed
2496
    default:
unknown's avatar
unknown committed
2497 2498 2499 2500 2501
      DBUG_ASSERT(1);
      return 0;
    }
  }
  Item_func_set_user_var e(user_var_name, it);
2502 2503 2504 2505
  /*
    Item_func_set_user_var can't substitute something else on its place =>
    0 can be passed as last argument (reference on item)
  */
unknown's avatar
unknown committed
2506
  e.fix_fields(thd, 0, 0);
2507
  e.update_hash(val, val_len, type, charset, DERIVATION_NONE);
unknown's avatar
unknown committed
2508 2509
  free_root(&thd->mem_root,0);

2510
  rli->inc_event_relay_log_pos(get_event_len());
unknown's avatar
unknown committed
2511 2512
  return 0;
}
unknown's avatar
unknown committed
2513
#endif /* !MYSQL_CLIENT */
2514 2515


unknown's avatar
unknown committed
2516
/**************************************************************************
2517
  Slave_log_event methods
unknown's avatar
unknown committed
2518
**************************************************************************/
unknown's avatar
unknown committed
2519

unknown's avatar
SCRUM  
unknown committed
2520
#ifdef HAVE_REPLICATION
2521 2522 2523 2524 2525 2526 2527 2528 2529 2530
#ifdef MYSQL_CLIENT
void Unknown_log_event::print(FILE* file, bool short_form, char* last_db)
{
  if (short_form)
    return;
  print_header(file);
  fputc('\n', file);
  fprintf(file, "# %s", "Unknown event\n");
}
#endif  
2531

2532
#ifndef MYSQL_CLIENT
2533
void Slave_log_event::pack_info(Protocol *protocol)
2534
{
unknown's avatar
unknown committed
2535
  char buf[256+HOSTNAME_LENGTH], *pos;
2536 2537 2538 2539 2540 2541 2542
  pos= strmov(buf, "host=");
  pos= strnmov(pos, master_host, HOSTNAME_LENGTH);
  pos= strmov(pos, ",port=");
  pos= int10_to_str((long) master_port, pos, 10);
  pos= strmov(pos, ",log=");
  pos= strmov(pos, master_log);
  pos= strmov(pos, ",pos=");
unknown's avatar
unknown committed
2543
  pos= longlong10_to_str(master_pos, pos, 10);
2544
  protocol->store(buf, pos-buf, &my_charset_bin);
2545
}
unknown's avatar
unknown committed
2546
#endif /* !MYSQL_CLIENT */
2547 2548 2549 2550


#ifndef MYSQL_CLIENT
Slave_log_event::Slave_log_event(THD* thd_arg,
unknown's avatar
unknown committed
2551 2552
				 struct st_relay_log_info* rli)
  :Log_event(thd_arg, 0, 0), mem_pool(0), master_host(0)
2553 2554 2555 2556 2557 2558 2559 2560 2561 2562
{
  DBUG_ENTER("Slave_log_event");
  if (!rli->inited)				// QQ When can this happen ?
    DBUG_VOID_RETURN;
  
  MASTER_INFO* mi = rli->mi;
  // TODO: re-write this better without holding both locks at the same time
  pthread_mutex_lock(&mi->data_lock);
  pthread_mutex_lock(&rli->data_lock);
  master_host_len = strlen(mi->host);
2563
  master_log_len = strlen(rli->group_master_log_name);
2564 2565 2566
  // on OOM, just do not initialize the structure and print the error
  if ((mem_pool = (char*)my_malloc(get_data_size() + 1,
				   MYF(MY_WME))))
2567
  {
2568 2569 2570
    master_host = mem_pool + SL_MASTER_HOST_OFFSET ;
    memcpy(master_host, mi->host, master_host_len + 1);
    master_log = master_host + master_host_len + 1;
2571
    memcpy(master_log, rli->group_master_log_name, master_log_len + 1);
2572
    master_port = mi->port;
2573
    master_pos = rli->group_master_log_pos;
2574 2575
    DBUG_PRINT("info", ("master_log: %s  pos: %d", master_log,
			(ulong) master_pos));
2576
  }
2577 2578 2579 2580 2581 2582
  else
    sql_print_error("Out of memory while recording slave event");
  pthread_mutex_unlock(&rli->data_lock);
  pthread_mutex_unlock(&mi->data_lock);
  DBUG_VOID_RETURN;
}
unknown's avatar
unknown committed
2583
#endif /* !MYSQL_CLIENT */
2584 2585 2586 2587 2588 2589 2590 2591 2592 2593 2594 2595 2596 2597 2598 2599


Slave_log_event::~Slave_log_event()
{
  my_free(mem_pool, MYF(MY_ALLOW_ZERO_PTR));
}


#ifdef MYSQL_CLIENT
void Slave_log_event::print(FILE* file, bool short_form, char* last_db)
{
  char llbuff[22];
  if (short_form)
    return;
  print_header(file);
  fputc('\n', file);
unknown's avatar
unknown committed
2600 2601
  fprintf(file, "\
Slave: master_host: '%s'  master_port: %d  master_log: '%s'  master_pos: %s\n",
2602 2603
	  master_host, master_port, master_log, llstr(master_pos, llbuff));
}
unknown's avatar
unknown committed
2604
#endif /* MYSQL_CLIENT */
2605 2606 2607 2608 2609 2610 2611 2612 2613 2614 2615 2616 2617 2618 2619 2620 2621 2622 2623 2624 2625 2626 2627 2628 2629 2630


int Slave_log_event::get_data_size()
{
  return master_host_len + master_log_len + 1 + SL_MASTER_HOST_OFFSET;
}


int Slave_log_event::write_data(IO_CACHE* file)
{
  int8store(mem_pool + SL_MASTER_POS_OFFSET, master_pos);
  int2store(mem_pool + SL_MASTER_PORT_OFFSET, master_port);
  // log and host are already there
  return my_b_safe_write(file, (byte*)mem_pool, get_data_size());
}


void Slave_log_event::init_from_mem_pool(int data_size)
{
  master_pos = uint8korr(mem_pool + SL_MASTER_POS_OFFSET);
  master_port = uint2korr(mem_pool + SL_MASTER_PORT_OFFSET);
  master_host = mem_pool + SL_MASTER_HOST_OFFSET;
  master_host_len = strlen(master_host);
  // safety
  master_log = master_host + master_host_len + 1;
  if (master_log > mem_pool + data_size)
2631
  {
2632 2633
    master_host = 0;
    return;
2634
  }
2635 2636
  master_log_len = strlen(master_log);
}
2637

2638 2639 2640 2641 2642 2643 2644 2645 2646 2647 2648 2649

Slave_log_event::Slave_log_event(const char* buf, int event_len)
  :Log_event(buf,0),mem_pool(0),master_host(0)
{
  event_len -= LOG_EVENT_HEADER_LEN;
  if (event_len < 0)
    return;
  if (!(mem_pool = (char*) my_malloc(event_len + 1, MYF(MY_WME))))
    return;
  memcpy(mem_pool, buf + LOG_EVENT_HEADER_LEN, event_len);
  mem_pool[event_len] = 0;
  init_from_mem_pool(event_len);
2650 2651
}

2652

2653 2654 2655 2656 2657 2658 2659
#ifndef MYSQL_CLIENT
int Slave_log_event::exec_event(struct st_relay_log_info* rli)
{
  if (mysql_bin_log.is_open())
    mysql_bin_log.write(this);
  return Log_event::exec_event(rli);
}
unknown's avatar
unknown committed
2660
#endif /* !MYSQL_CLIENT */
2661 2662


unknown's avatar
unknown committed
2663
/**************************************************************************
unknown's avatar
unknown committed
2664
	Stop_log_event methods
unknown's avatar
unknown committed
2665
**************************************************************************/
2666

unknown's avatar
unknown committed
2667
/*
2668
  Stop_log_event::print()
2669
*/
2670 2671 2672 2673 2674 2675 2676 2677 2678 2679

#ifdef MYSQL_CLIENT
void Stop_log_event::print(FILE* file, bool short_form, char* last_db)
{
  if (short_form)
    return;

  print_header(file);
  fprintf(file, "\tStop\n");
  fflush(file);
2680
}
unknown's avatar
unknown committed
2681
#endif /* MYSQL_CLIENT */
2682

2683

2684
/*
2685
  Stop_log_event::exec_event()
2686

2687
  The master stopped. 
unknown's avatar
unknown committed
2688 2689
  We used to clean up all temporary tables but this is useless as, as the
  master has shut down properly, it has written all DROP TEMPORARY TABLE and DO
2690 2691 2692 2693 2694 2695
  RELEASE_LOCK (prepared statements' deletion is TODO).
  We used to clean up slave_load_tmpdir, but this is useless as it has been
  cleared at the end of LOAD DATA INFILE.
  So we have nothing to do here.
  The place were we must do this cleaning is in Start_log_event::exec_event(),
  not here. Because if we come here, the master was sane.
2696 2697
*/

2698
#ifndef MYSQL_CLIENT
2699
int Stop_log_event::exec_event(struct st_relay_log_info* rli)
2700
{
unknown's avatar
unknown committed
2701 2702
  /*
    We do not want to update master_log pos because we get a rotate event
2703
    before stop, so by now group_master_log_name is set to the next log.
2704
    If we updated it, we will have incorrect master coordinates and this
unknown's avatar
unknown committed
2705
    could give false triggers in MASTER_POS_WAIT() that we have reached
2706
    the target position when in fact we have not.
unknown's avatar
unknown committed
2707
  */
unknown's avatar
unknown committed
2708
  rli->inc_group_relay_log_pos(get_event_len(), 0);
2709
  flush_relay_log_info(rli);
2710 2711
  return 0;
}
unknown's avatar
unknown committed
2712
#endif /* !MYSQL_CLIENT */
unknown's avatar
SCRUM  
unknown committed
2713
#endif /* HAVE_REPLICATION */
2714

2715

unknown's avatar
unknown committed
2716
/**************************************************************************
unknown's avatar
unknown committed
2717
	Create_file_log_event methods
unknown's avatar
unknown committed
2718
**************************************************************************/
2719 2720

/*
2721
  Create_file_log_event ctor
unknown's avatar
unknown committed
2722
*/
2723 2724

#ifndef MYSQL_CLIENT
unknown's avatar
unknown committed
2725 2726 2727 2728 2729 2730 2731
Create_file_log_event::
Create_file_log_event(THD* thd_arg, sql_exchange* ex,
		      const char* db_arg, const char* table_name_arg,
		      List<Item>& fields_arg, enum enum_duplicates handle_dup,
		      char* block_arg, uint block_len_arg, bool using_trans)
  :Load_log_event(thd_arg,ex,db_arg,table_name_arg,fields_arg,handle_dup,
		  using_trans),
2732
   fake_base(0), block(block_arg), event_buf(0), block_len(block_len_arg),
2733
   file_id(thd_arg->file_id = mysql_bin_log.next_file_id())
2734
{
unknown's avatar
unknown committed
2735
  DBUG_ENTER("Create_file_log_event");
2736
  sql_ex.force_new_format();
unknown's avatar
unknown committed
2737
  DBUG_VOID_RETURN;
2738
}
unknown's avatar
unknown committed
2739
#endif /* !MYSQL_CLIENT */
2740

2741

unknown's avatar
unknown committed
2742
/*
2743
  Create_file_log_event::write_data_body()
unknown's avatar
unknown committed
2744
*/
2745 2746 2747 2748 2749 2750 2751 2752

int Create_file_log_event::write_data_body(IO_CACHE* file)
{
  int res;
  if ((res = Load_log_event::write_data_body(file)) || fake_base)
    return res;
  return (my_b_safe_write(file, (byte*) "", 1) ||
	  my_b_safe_write(file, (byte*) block, block_len));
2753 2754
}

2755

unknown's avatar
unknown committed
2756
/*
2757
  Create_file_log_event::write_data_header()
unknown's avatar
unknown committed
2758
*/
unknown's avatar
unknown committed
2759

2760
int Create_file_log_event::write_data_header(IO_CACHE* file)
2761
{
2762 2763 2764 2765 2766 2767 2768 2769 2770
  int res;
  if ((res = Load_log_event::write_data_header(file)) || fake_base)
    return res;
  byte buf[CREATE_FILE_HEADER_LEN];
  int4store(buf + CF_FILE_ID_OFFSET, file_id);
  return my_b_safe_write(file, buf, CREATE_FILE_HEADER_LEN);
}


unknown's avatar
unknown committed
2771
/*
2772
  Create_file_log_event::write_base()
unknown's avatar
unknown committed
2773
*/
2774 2775 2776 2777 2778 2779 2780 2781 2782 2783 2784

int Create_file_log_event::write_base(IO_CACHE* file)
{
  int res;
  fake_base = 1; // pretend we are Load event
  res = write(file);
  fake_base = 0;
  return res;
}


unknown's avatar
unknown committed
2785
/*
2786
  Create_file_log_event ctor
unknown's avatar
unknown committed
2787
*/
2788 2789 2790 2791 2792 2793

Create_file_log_event::Create_file_log_event(const char* buf, int len,
					     bool old_format)
  :Load_log_event(buf,0,old_format),fake_base(0),block(0),inited_from_old(0)
{
  int block_offset;
unknown's avatar
unknown committed
2794 2795
  DBUG_ENTER("Create_file_log_event");

2796
  /*
2797 2798
    We must make copy of 'buf' as this event may have to live over a
    rotate log entry when used in mysqlbinlog
2799
  */
unknown's avatar
unknown committed
2800
  if (!(event_buf= my_memdup((byte*) buf, len, MYF(MY_WME))) ||
2801
      (copy_log_event(event_buf, len, old_format)))
unknown's avatar
unknown committed
2802
    DBUG_VOID_RETURN;
2803

2804 2805 2806 2807 2808 2809 2810 2811 2812 2813 2814 2815 2816 2817 2818 2819
  if (!old_format)
  {
    file_id = uint4korr(buf + LOG_EVENT_HEADER_LEN +
			+ LOAD_HEADER_LEN + CF_FILE_ID_OFFSET);
    // + 1 for \0 terminating fname  
    block_offset = (LOG_EVENT_HEADER_LEN + Load_log_event::get_data_size() +
		    CREATE_FILE_HEADER_LEN + 1);
    if (len < block_offset)
      return;
    block = (char*)buf + block_offset;
    block_len = len - block_offset;
  }
  else
  {
    sql_ex.force_new_format();
    inited_from_old = 1;
2820
  }
unknown's avatar
unknown committed
2821
  DBUG_VOID_RETURN;
2822 2823
}

2824

unknown's avatar
unknown committed
2825
/*
2826
  Create_file_log_event::print()
unknown's avatar
unknown committed
2827
*/
2828 2829

#ifdef MYSQL_CLIENT
2830 2831
void Create_file_log_event::print(FILE* file, bool short_form, 
				  char* last_db, bool enable_local)
unknown's avatar
unknown committed
2832
{
2833
  if (short_form)
2834 2835 2836
  {
    if (enable_local && check_fname_outside_temp_buf())
      Load_log_event::print(file, 1, last_db);
2837
    return;
2838 2839 2840 2841
  }

  if (enable_local)
  {
unknown's avatar
unknown committed
2842 2843 2844 2845 2846 2847
    Load_log_event::print(file, 1, last_db, !check_fname_outside_temp_buf());
    /* 
       That one is for "file_id: etc" below: in mysqlbinlog we want the #, in
       SHOW BINLOG EVENTS we don't.
    */
    fprintf(file, "#"); 
2848 2849
  }

2850
  fprintf(file, " file_id: %d  block_len: %d\n", file_id, block_len);
unknown's avatar
unknown committed
2851
}
2852

unknown's avatar
unknown committed
2853

2854 2855 2856 2857 2858
void Create_file_log_event::print(FILE* file, bool short_form,
				  char* last_db)
{
  print(file,short_form,last_db,0);
}
unknown's avatar
unknown committed
2859
#endif /* MYSQL_CLIENT */
unknown's avatar
unknown committed
2860

2861

unknown's avatar
unknown committed
2862
/*
2863
  Create_file_log_event::pack_info()
unknown's avatar
unknown committed
2864
*/
2865

unknown's avatar
SCRUM  
unknown committed
2866
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
2867
void Create_file_log_event::pack_info(Protocol *protocol)
2868
{
2869 2870 2871
  char buf[NAME_LEN*2 + 30 + 21*2], *pos;
  pos= strmov(buf, "db=");
  memcpy(pos, db, db_len);
unknown's avatar
unknown committed
2872
  pos= strmov(pos + db_len, ";table=");
2873
  memcpy(pos, table_name, table_name_len);
unknown's avatar
unknown committed
2874
  pos= strmov(pos + table_name_len, ";file_id=");
2875 2876 2877
  pos= int10_to_str((long) file_id, pos, 10);
  pos= strmov(pos, ";block_len=");
  pos= int10_to_str((long) block_len, pos, 10);
unknown's avatar
unknown committed
2878
  protocol->store(buf, (uint) (pos-buf), &my_charset_bin);
2879
}
unknown's avatar
unknown committed
2880
#endif /* defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT) */
2881 2882


unknown's avatar
unknown committed
2883
/*
2884
  Create_file_log_event::exec_event()
unknown's avatar
unknown committed
2885
*/
2886

unknown's avatar
SCRUM  
unknown committed
2887
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
2888
int Create_file_log_event::exec_event(struct st_relay_log_info* rli)
2889
{
2890
  char proc_info[17+FN_REFLEN+10], *fname_buf= proc_info+17;
unknown's avatar
unknown committed
2891
  char *p;
2892 2893 2894
  int fd = -1;
  IO_CACHE file;
  int error = 1;
unknown's avatar
unknown committed
2895

2896
  bzero((char*)&file, sizeof(file));
unknown's avatar
unknown committed
2897 2898
  p = slave_load_file_stem(fname_buf, file_id, server_id);
  strmov(p, ".info");			// strmov takes less code than memcpy
2899 2900
  strnmov(proc_info, "Making temp file ", 17); // no end 0
  thd->proc_info= proc_info;
2901 2902 2903 2904 2905
  if ((fd = my_open(fname_buf, O_WRONLY|O_CREAT|O_BINARY|O_TRUNC,
		    MYF(MY_WME))) < 0 ||
      init_io_cache(&file, fd, IO_SIZE, WRITE_CACHE, (my_off_t)0, 0,
		    MYF(MY_WME|MY_NABP)))
  {
unknown's avatar
unknown committed
2906
    slave_print_error(rli,my_errno, "Error in Create_file event: could not open file '%s'", fname_buf);
2907 2908 2909 2910
    goto err;
  }
  
  // a trick to avoid allocating another buffer
unknown's avatar
unknown committed
2911
  strmov(p, ".data");
2912 2913 2914 2915
  fname = fname_buf;
  fname_len = (uint)(p-fname) + 5;
  if (write_base(&file))
  {
unknown's avatar
unknown committed
2916
    strmov(p, ".info"); // to have it right in the error message
unknown's avatar
unknown committed
2917 2918 2919
    slave_print_error(rli,my_errno,
		      "Error in Create_file event: could not write to file '%s'",
		      fname_buf);
2920 2921 2922 2923 2924 2925 2926 2927 2928
    goto err;
  }
  end_io_cache(&file);
  my_close(fd, MYF(0));
  
  // fname_buf now already has .data, not .info, because we did our trick
  if ((fd = my_open(fname_buf, O_WRONLY|O_CREAT|O_BINARY|O_TRUNC,
		    MYF(MY_WME))) < 0)
  {
unknown's avatar
unknown committed
2929
    slave_print_error(rli,my_errno, "Error in Create_file event: could not open file '%s'", fname_buf);
2930 2931
    goto err;
  }
unknown's avatar
unknown committed
2932
  if (my_write(fd, (byte*) block, block_len, MYF(MY_WME+MY_NABP)))
2933
  {
unknown's avatar
unknown committed
2934
    slave_print_error(rli,my_errno, "Error in Create_file event: write to '%s' failed", fname_buf);
2935 2936
    goto err;
  }
2937 2938
  error=0;					// Everything is ok

2939 2940 2941 2942 2943
err:
  if (error)
    end_io_cache(&file);
  if (fd >= 0)
    my_close(fd, MYF(0));
unknown's avatar
unknown committed
2944
  thd->proc_info= 0;
2945
  return error ? 1 : Log_event::exec_event(rli);
2946
}
unknown's avatar
unknown committed
2947
#endif /* defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT) */
2948

2949

unknown's avatar
unknown committed
2950
/**************************************************************************
unknown's avatar
unknown committed
2951
	Append_block_log_event methods
unknown's avatar
unknown committed
2952
**************************************************************************/
2953

unknown's avatar
unknown committed
2954
/*
2955
  Append_block_log_event ctor
unknown's avatar
unknown committed
2956
*/
2957 2958

#ifndef MYSQL_CLIENT  
unknown's avatar
unknown committed
2959 2960
Append_block_log_event::Append_block_log_event(THD* thd_arg, const char* db_arg,
					       char* block_arg,
unknown's avatar
unknown committed
2961 2962 2963
					       uint block_len_arg,
					       bool using_trans)
  :Log_event(thd_arg,0, using_trans), block(block_arg),
unknown's avatar
unknown committed
2964
   block_len(block_len_arg), file_id(thd_arg->file_id), db(db_arg)
2965 2966
{
}
unknown's avatar
unknown committed
2967
#endif
2968 2969


unknown's avatar
unknown committed
2970
/*
2971
  Append_block_log_event ctor
unknown's avatar
unknown committed
2972
*/
2973 2974 2975 2976

Append_block_log_event::Append_block_log_event(const char* buf, int len)
  :Log_event(buf, 0),block(0)
{
unknown's avatar
unknown committed
2977
  DBUG_ENTER("Append_block_log_event");
2978
  if ((uint)len < APPEND_BLOCK_EVENT_OVERHEAD)
unknown's avatar
unknown committed
2979
    DBUG_VOID_RETURN;
2980 2981 2982
  file_id = uint4korr(buf + LOG_EVENT_HEADER_LEN + AB_FILE_ID_OFFSET);
  block = (char*)buf + APPEND_BLOCK_EVENT_OVERHEAD;
  block_len = len - APPEND_BLOCK_EVENT_OVERHEAD;
unknown's avatar
unknown committed
2983
  DBUG_VOID_RETURN;
2984 2985 2986
}


unknown's avatar
unknown committed
2987
/*
2988
  Append_block_log_event::write_data()
unknown's avatar
unknown committed
2989
*/
2990 2991 2992 2993 2994 2995 2996 2997 2998 2999

int Append_block_log_event::write_data(IO_CACHE* file)
{
  byte buf[APPEND_BLOCK_HEADER_LEN];
  int4store(buf + AB_FILE_ID_OFFSET, file_id);
  return (my_b_safe_write(file, buf, APPEND_BLOCK_HEADER_LEN) ||
	  my_b_safe_write(file, (byte*) block, block_len));
}


unknown's avatar
unknown committed
3000
/*
3001
  Append_block_log_event::print()
unknown's avatar
unknown committed
3002
*/
3003 3004 3005 3006 3007 3008 3009 3010 3011 3012 3013 3014

#ifdef MYSQL_CLIENT  
void Append_block_log_event::print(FILE* file, bool short_form,
				   char* last_db)
{
  if (short_form)
    return;
  print_header(file);
  fputc('\n', file);
  fprintf(file, "#Append_block: file_id: %d  block_len: %d\n",
	  file_id, block_len);
}
unknown's avatar
unknown committed
3015
#endif /* MYSQL_CLIENT */
3016 3017


unknown's avatar
unknown committed
3018
/*
3019
  Append_block_log_event::pack_info()
unknown's avatar
unknown committed
3020
*/
3021

unknown's avatar
SCRUM  
unknown committed
3022
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
3023
void Append_block_log_event::pack_info(Protocol *protocol)
3024 3025 3026 3027 3028 3029
{
  char buf[256];
  uint length;
  length= (uint) my_sprintf(buf,
			    (buf, ";file_id=%u;block_len=%u", file_id,
			     block_len));
unknown's avatar
unknown committed
3030
  protocol->store(buf, length, &my_charset_bin);
3031
}
unknown's avatar
unknown committed
3032
#endif /* defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT) */
3033 3034


unknown's avatar
unknown committed
3035
/*
3036
  Append_block_log_event::exec_event()
unknown's avatar
unknown committed
3037
*/
3038

unknown's avatar
SCRUM  
unknown committed
3039
#if defined( HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
3040
int Append_block_log_event::exec_event(struct st_relay_log_info* rli)
3041
{
3042
  char proc_info[17+FN_REFLEN+10], *fname= proc_info+17;
3043 3044
  char *p= slave_load_file_stem(fname, file_id, server_id);
  int fd;
3045
  int error = 1;
unknown's avatar
unknown committed
3046
  DBUG_ENTER("Append_block_log_event::exec_event");
3047

3048
  memcpy(p, ".data", 6);
3049 3050
  strnmov(proc_info, "Making temp file ", 17); // no end 0
  thd->proc_info= proc_info;
3051 3052
  if ((fd = my_open(fname, O_WRONLY|O_APPEND|O_BINARY, MYF(MY_WME))) < 0)
  {
unknown's avatar
unknown committed
3053
    slave_print_error(rli,my_errno, "Error in Append_block event: could not open file '%s'", fname);
3054 3055
    goto err;
  }
unknown's avatar
unknown committed
3056
  if (my_write(fd, (byte*) block, block_len, MYF(MY_WME+MY_NABP)))
3057
  {
unknown's avatar
unknown committed
3058
    slave_print_error(rli,my_errno, "Error in Append_block event: write to '%s' failed", fname);
3059 3060 3061
    goto err;
  }
  error=0;
3062

3063 3064 3065
err:
  if (fd >= 0)
    my_close(fd, MYF(0));
3066
  thd->proc_info= 0;
unknown's avatar
unknown committed
3067
  DBUG_RETURN(error ? error : Log_event::exec_event(rli));
3068
}
unknown's avatar
SCRUM  
unknown committed
3069
#endif
3070 3071


unknown's avatar
unknown committed
3072
/**************************************************************************
unknown's avatar
unknown committed
3073
	Delete_file_log_event methods
unknown's avatar
unknown committed
3074
**************************************************************************/
3075

unknown's avatar
unknown committed
3076
/*
3077
  Delete_file_log_event ctor
unknown's avatar
unknown committed
3078
*/
3079 3080

#ifndef MYSQL_CLIENT
unknown's avatar
unknown committed
3081 3082 3083
Delete_file_log_event::Delete_file_log_event(THD *thd_arg, const char* db_arg,
					     bool using_trans)
  :Log_event(thd_arg, 0, using_trans), file_id(thd_arg->file_id), db(db_arg)
3084 3085
{
}
unknown's avatar
unknown committed
3086
#endif
3087

unknown's avatar
unknown committed
3088
/*
3089
  Delete_file_log_event ctor
unknown's avatar
unknown committed
3090
*/
3091 3092 3093 3094 3095 3096 3097 3098 3099 3100

Delete_file_log_event::Delete_file_log_event(const char* buf, int len)
  :Log_event(buf, 0),file_id(0)
{
  if ((uint)len < DELETE_FILE_EVENT_OVERHEAD)
    return;
  file_id = uint4korr(buf + LOG_EVENT_HEADER_LEN + AB_FILE_ID_OFFSET);
}


unknown's avatar
unknown committed
3101
/*
3102
  Delete_file_log_event::write_data()
unknown's avatar
unknown committed
3103
*/
3104 3105 3106 3107 3108 3109 3110 3111 3112

int Delete_file_log_event::write_data(IO_CACHE* file)
{
 byte buf[DELETE_FILE_HEADER_LEN];
 int4store(buf + DF_FILE_ID_OFFSET, file_id);
 return my_b_safe_write(file, buf, DELETE_FILE_HEADER_LEN);
}


unknown's avatar
unknown committed
3113
/*
3114
  Delete_file_log_event::print()
unknown's avatar
unknown committed
3115
*/
3116 3117 3118 3119 3120 3121 3122 3123 3124 3125 3126

#ifdef MYSQL_CLIENT  
void Delete_file_log_event::print(FILE* file, bool short_form,
				  char* last_db)
{
  if (short_form)
    return;
  print_header(file);
  fputc('\n', file);
  fprintf(file, "#Delete_file: file_id=%u\n", file_id);
}
unknown's avatar
unknown committed
3127
#endif /* MYSQL_CLIENT */
3128

unknown's avatar
unknown committed
3129
/*
3130
  Delete_file_log_event::pack_info()
unknown's avatar
unknown committed
3131
*/
3132

unknown's avatar
SCRUM  
unknown committed
3133
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
3134
void Delete_file_log_event::pack_info(Protocol *protocol)
3135 3136 3137 3138
{
  char buf[64];
  uint length;
  length= (uint) my_sprintf(buf, (buf, ";file_id=%u", (uint) file_id));
3139
  protocol->store(buf, (int32) length, &my_charset_bin);
3140
}
unknown's avatar
SCRUM  
unknown committed
3141
#endif
3142

unknown's avatar
unknown committed
3143
/*
3144
  Delete_file_log_event::exec_event()
unknown's avatar
unknown committed
3145
*/
3146

unknown's avatar
SCRUM  
unknown committed
3147
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
3148 3149 3150 3151 3152 3153 3154 3155 3156 3157
int Delete_file_log_event::exec_event(struct st_relay_log_info* rli)
{
  char fname[FN_REFLEN+10];
  char *p= slave_load_file_stem(fname, file_id, server_id);
  memcpy(p, ".data", 6);
  (void) my_delete(fname, MYF(MY_WME));
  memcpy(p, ".info", 6);
  (void) my_delete(fname, MYF(MY_WME));
  return Log_event::exec_event(rli);
}
unknown's avatar
unknown committed
3158
#endif /* defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT) */
3159 3160


unknown's avatar
unknown committed
3161
/**************************************************************************
unknown's avatar
unknown committed
3162
	Execute_load_log_event methods
unknown's avatar
unknown committed
3163
**************************************************************************/
3164

unknown's avatar
unknown committed
3165
/*
3166
  Execute_load_log_event ctor
unknown's avatar
unknown committed
3167
*/
3168 3169

#ifndef MYSQL_CLIENT  
unknown's avatar
unknown committed
3170 3171 3172
Execute_load_log_event::Execute_load_log_event(THD *thd_arg, const char* db_arg,
					       bool using_trans)
  :Log_event(thd_arg, 0, using_trans), file_id(thd_arg->file_id), db(db_arg)
3173 3174
{
}
unknown's avatar
unknown committed
3175
#endif
3176 3177
  

unknown's avatar
unknown committed
3178
/*
3179
  Execute_load_log_event ctor
unknown's avatar
unknown committed
3180
*/
3181

unknown's avatar
unknown committed
3182 3183
Execute_load_log_event::Execute_load_log_event(const char* buf, int len)
  :Log_event(buf, 0), file_id(0)
3184 3185 3186 3187 3188 3189 3190
{
  if ((uint)len < EXEC_LOAD_EVENT_OVERHEAD)
    return;
  file_id = uint4korr(buf + LOG_EVENT_HEADER_LEN + EL_FILE_ID_OFFSET);
}


unknown's avatar
unknown committed
3191
/*
3192
  Execute_load_log_event::write_data()
unknown's avatar
unknown committed
3193
*/
3194 3195 3196 3197 3198 3199 3200 3201 3202

int Execute_load_log_event::write_data(IO_CACHE* file)
{
  byte buf[EXEC_LOAD_HEADER_LEN];
  int4store(buf + EL_FILE_ID_OFFSET, file_id);
  return my_b_safe_write(file, buf, EXEC_LOAD_HEADER_LEN);
}


unknown's avatar
unknown committed
3203
/*
3204
  Execute_load_log_event::print()
unknown's avatar
unknown committed
3205
*/
3206 3207 3208 3209 3210 3211 3212 3213 3214 3215 3216 3217

#ifdef MYSQL_CLIENT  
void Execute_load_log_event::print(FILE* file, bool short_form,
				   char* last_db)
{
  if (short_form)
    return;
  print_header(file);
  fputc('\n', file);
  fprintf(file, "#Exec_load: file_id=%d\n",
	  file_id);
}
unknown's avatar
unknown committed
3218
#endif
3219

unknown's avatar
unknown committed
3220
/*
3221
  Execute_load_log_event::pack_info()
unknown's avatar
unknown committed
3222
*/
3223

unknown's avatar
SCRUM  
unknown committed
3224
#if defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT)
3225
void Execute_load_log_event::pack_info(Protocol *protocol)
3226 3227 3228 3229
{
  char buf[64];
  uint length;
  length= (uint) my_sprintf(buf, (buf, ";file_id=%u", (uint) file_id));
3230
  protocol->store(buf, (int32) length, &my_charset_bin);
3231 3232 3233
}


unknown's avatar
unknown committed
3234
/*
3235
  Execute_load_log_event::exec_event()
unknown's avatar
unknown committed
3236
*/
unknown's avatar
SCRUM  
unknown committed
3237

3238
int Execute_load_log_event::exec_event(struct st_relay_log_info* rli)
3239 3240
{
  char fname[FN_REFLEN+10];
3241 3242
  char *p= slave_load_file_stem(fname, file_id, server_id);
  int fd;
3243 3244 3245
  int error = 1;
  IO_CACHE file;
  Load_log_event* lev = 0;
3246

3247 3248 3249 3250 3251
  memcpy(p, ".info", 6);
  if ((fd = my_open(fname, O_RDONLY|O_BINARY, MYF(MY_WME))) < 0 ||
      init_io_cache(&file, fd, IO_SIZE, READ_CACHE, (my_off_t)0, 0,
		    MYF(MY_WME|MY_NABP)))
  {
unknown's avatar
unknown committed
3252
    slave_print_error(rli,my_errno, "Error in Exec_load event: could not open file '%s'", fname);
3253 3254
    goto err;
  }
3255 3256
  if (!(lev = (Load_log_event*)Log_event::read_log_event(&file,
							 (pthread_mutex_t*)0,
3257 3258
							 (bool)0)) ||
      lev->get_type_code() != NEW_LOAD_EVENT)
3259
  {
unknown's avatar
unknown committed
3260
    slave_print_error(rli,0, "Error in Exec_load event: file '%s' appears corrupted", fname);
3261 3262
    goto err;
  }
unknown's avatar
unknown committed
3263

3264
  lev->thd = thd;
3265 3266
  /*
    lev->exec_event should use rli only for errors
unknown's avatar
unknown committed
3267 3268 3269
    i.e. should not advance rli's position.
    lev->exec_event is the place where the table is loaded (it calls
    mysql_load()).
3270
  */
unknown's avatar
unknown committed
3271 3272

#if MYSQL_VERSION_ID < 40100
3273 3274
    rli->future_master_log_pos= log_pos + get_event_len() -
      (rli->mi->old_format ? (LOG_EVENT_HEADER_LEN - OLD_HEADER_LEN) : 0);
unknown's avatar
unknown committed
3275
#elif MYSQL_VERSION_ID < 50000
3276 3277
    rli->future_group_master_log_pos= log_pos + get_event_len() -
      (rli->mi->old_format ? (LOG_EVENT_HEADER_LEN - OLD_HEADER_LEN) : 0);
unknown's avatar
unknown committed
3278 3279 3280
#else
    rli->future_group_master_log_pos= log_pos;
#endif
3281
  if (lev->exec_event(0,rli,1)) 
3282
  {
3283 3284 3285 3286 3287 3288 3289 3290 3291
    /*
      We want to indicate the name of the file that could not be loaded
      (SQL_LOADxxx).
      But as we are here we are sure the error is in rli->last_slave_error and
      rli->last_slave_errno (example of error: duplicate entry for key), so we
      don't want to overwrite it with the filename.
      What we want instead is add the filename to the current error message.
    */
    char *tmp= my_strdup(rli->last_slave_error,MYF(MY_WME));
3292 3293 3294 3295 3296 3297 3298 3299
    if (tmp)
    {
      slave_print_error(rli,
			rli->last_slave_errno, /* ok to re-use error code */
			"%s. Failed executing load from '%s'", 
			tmp, fname);
      my_free(tmp,MYF(0));
    }
3300 3301
    goto err;
  }
unknown's avatar
unknown committed
3302 3303 3304 3305 3306 3307 3308 3309 3310 3311
  /*
    We have an open file descriptor to the .info file; we need to close it
    or Windows will refuse to delete the file in my_delete().
  */
  if (fd >= 0)
  {
    my_close(fd, MYF(0));
    end_io_cache(&file);
    fd= -1;
  }
3312
  (void) my_delete(fname, MYF(MY_WME));
3313
  memcpy(p, ".data", 6);
3314
  (void) my_delete(fname, MYF(MY_WME));
3315
  error = 0;
3316

3317 3318 3319
err:
  delete lev;
  if (fd >= 0)
3320
  {
3321
    my_close(fd, MYF(0));
3322 3323
    end_io_cache(&file);
  }
3324
  return error ? error : Log_event::exec_event(rli);
3325
}
unknown's avatar
SCRUM  
unknown committed
3326

unknown's avatar
unknown committed
3327
#endif /* defined(HAVE_REPLICATION) && !defined(MYSQL_CLIENT) */
3328 3329


unknown's avatar
unknown committed
3330
/**************************************************************************
unknown's avatar
unknown committed
3331
	sql_ex_info methods
unknown's avatar
unknown committed
3332
**************************************************************************/
3333

unknown's avatar
unknown committed
3334
/*
3335
  sql_ex_info::write_data()
unknown's avatar
unknown committed
3336
*/
3337 3338 3339 3340 3341 3342 3343 3344 3345 3346 3347 3348 3349 3350 3351 3352 3353 3354 3355 3356 3357 3358 3359 3360 3361 3362 3363

int sql_ex_info::write_data(IO_CACHE* file)
{
  if (new_format())
  {
    return (write_str(file, field_term, field_term_len) ||
	    write_str(file, enclosed,   enclosed_len) ||
	    write_str(file, line_term,  line_term_len) ||
	    write_str(file, line_start, line_start_len) ||
	    write_str(file, escaped,    escaped_len) ||
	    my_b_safe_write(file,(byte*) &opt_flags,1));
  }
  else
  {
    old_sql_ex old_ex;
    old_ex.field_term= *field_term;
    old_ex.enclosed=   *enclosed;
    old_ex.line_term=  *line_term;
    old_ex.line_start= *line_start;
    old_ex.escaped=    *escaped;
    old_ex.opt_flags=  opt_flags;
    old_ex.empty_flags=empty_flags;
    return my_b_safe_write(file, (byte*) &old_ex, sizeof(old_ex));
  }
}


unknown's avatar
unknown committed
3364
/*
3365
  sql_ex_info::init()
unknown's avatar
unknown committed
3366
*/
3367 3368 3369 3370 3371 3372 3373 3374 3375 3376 3377 3378 3379 3380 3381 3382 3383 3384 3385 3386 3387 3388 3389 3390 3391 3392 3393 3394 3395 3396 3397 3398 3399 3400 3401 3402 3403 3404 3405 3406 3407 3408 3409 3410 3411

char* sql_ex_info::init(char* buf,char* buf_end,bool use_new_format)
{
  cached_new_format = use_new_format;
  if (use_new_format)
  {
    empty_flags=0;
    /*
      The code below assumes that buf will not disappear from
      under our feet during the lifetime of the event. This assumption
      holds true in the slave thread if the log is in new format, but is not
      the case when we have old format because we will be reusing net buffer
      to read the actual file before we write out the Create_file event.
    */
    if (read_str(buf, buf_end, field_term, field_term_len) ||
	read_str(buf, buf_end, enclosed,   enclosed_len) ||
	read_str(buf, buf_end, line_term,  line_term_len) ||
	read_str(buf, buf_end, line_start, line_start_len) ||
	read_str(buf, buf_end, escaped,	   escaped_len))
      return 0;
    opt_flags = *buf++;
  }
  else
  {
    field_term_len= enclosed_len= line_term_len= line_start_len= escaped_len=1;
    field_term = buf++;			// Use first byte in string
    enclosed=	 buf++;
    line_term=   buf++;
    line_start=  buf++;
    escaped=     buf++;
    opt_flags =  *buf++;
    empty_flags= *buf++;
    if (empty_flags & FIELD_TERM_EMPTY)
      field_term_len=0;
    if (empty_flags & ENCLOSED_EMPTY)
      enclosed_len=0;
    if (empty_flags & LINE_TERM_EMPTY)
      line_term_len=0;
    if (empty_flags & LINE_START_EMPTY)
      line_start_len=0;
    if (empty_flags & ESCAPED_EMPTY)
      escaped_len=0;
  }
  return buf;
}