|
|
|
@ -25,26 +25,23 @@ import ( |
|
|
|
|
|
|
|
|
|
// Read packet to buffer 'data'
|
|
|
|
|
func (mc *mysqlConn) readPacket() ([]byte, error) { |
|
|
|
|
var payload []byte |
|
|
|
|
var prevData []byte |
|
|
|
|
for { |
|
|
|
|
// Read packet header
|
|
|
|
|
// read packet header
|
|
|
|
|
data, err := mc.buf.readNext(4) |
|
|
|
|
if err != nil { |
|
|
|
|
if cerr := mc.canceled.Value(); cerr != nil { |
|
|
|
|
return nil, cerr |
|
|
|
|
} |
|
|
|
|
errLog.Print(err) |
|
|
|
|
mc.Close() |
|
|
|
|
return nil, driver.ErrBadConn |
|
|
|
|
return nil, ErrInvalidConn |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Packet Length [24 bit]
|
|
|
|
|
// packet length [24 bit]
|
|
|
|
|
pktLen := int(uint32(data[0]) | uint32(data[1])<<8 | uint32(data[2])<<16) |
|
|
|
|
|
|
|
|
|
if pktLen < 1 { |
|
|
|
|
errLog.Print(ErrMalformPkt) |
|
|
|
|
mc.Close() |
|
|
|
|
return nil, driver.ErrBadConn |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Check Packet Sync [8 bit]
|
|
|
|
|
// check packet sync [8 bit]
|
|
|
|
|
if data[3] != mc.sequence { |
|
|
|
|
if data[3] > mc.sequence { |
|
|
|
|
return nil, ErrPktSyncMul |
|
|
|
@ -53,26 +50,41 @@ func (mc *mysqlConn) readPacket() ([]byte, error) { |
|
|
|
|
} |
|
|
|
|
mc.sequence++ |
|
|
|
|
|
|
|
|
|
// Read packet body [pktLen bytes]
|
|
|
|
|
// packets with length 0 terminate a previous packet which is a
|
|
|
|
|
// multiple of (2^24)−1 bytes long
|
|
|
|
|
if pktLen == 0 { |
|
|
|
|
// there was no previous packet
|
|
|
|
|
if prevData == nil { |
|
|
|
|
errLog.Print(ErrMalformPkt) |
|
|
|
|
mc.Close() |
|
|
|
|
return nil, ErrInvalidConn |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
return prevData, nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// read packet body [pktLen bytes]
|
|
|
|
|
data, err = mc.buf.readNext(pktLen) |
|
|
|
|
if err != nil { |
|
|
|
|
if cerr := mc.canceled.Value(); cerr != nil { |
|
|
|
|
return nil, cerr |
|
|
|
|
} |
|
|
|
|
errLog.Print(err) |
|
|
|
|
mc.Close() |
|
|
|
|
return nil, driver.ErrBadConn |
|
|
|
|
return nil, ErrInvalidConn |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
isLastPacket := (pktLen < maxPacketSize) |
|
|
|
|
|
|
|
|
|
// Zero allocations for non-splitting packets
|
|
|
|
|
if isLastPacket && payload == nil { |
|
|
|
|
// return data if this was the last packet
|
|
|
|
|
if pktLen < maxPacketSize { |
|
|
|
|
// zero allocations for non-split packets
|
|
|
|
|
if prevData == nil { |
|
|
|
|
return data, nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
payload = append(payload, data...) |
|
|
|
|
|
|
|
|
|
if isLastPacket { |
|
|
|
|
return payload, nil |
|
|
|
|
return append(prevData, data...), nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
prevData = append(prevData, data...) |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
@ -119,33 +131,47 @@ func (mc *mysqlConn) writePacket(data []byte) error { |
|
|
|
|
|
|
|
|
|
// Handle error
|
|
|
|
|
if err == nil { // n != len(data)
|
|
|
|
|
mc.cleanup() |
|
|
|
|
errLog.Print(ErrMalformPkt) |
|
|
|
|
} else { |
|
|
|
|
if cerr := mc.canceled.Value(); cerr != nil { |
|
|
|
|
return cerr |
|
|
|
|
} |
|
|
|
|
if n == 0 && pktLen == len(data)-4 { |
|
|
|
|
// only for the first loop iteration when nothing was written yet
|
|
|
|
|
return errBadConnNoWrite |
|
|
|
|
} |
|
|
|
|
mc.cleanup() |
|
|
|
|
errLog.Print(err) |
|
|
|
|
} |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
return ErrInvalidConn |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/****************************************************************************** |
|
|
|
|
* Initialisation Process * |
|
|
|
|
* Initialization Process * |
|
|
|
|
******************************************************************************/ |
|
|
|
|
|
|
|
|
|
// Handshake Initialization Packet
|
|
|
|
|
// http://dev.mysql.com/doc/internals/en/connection-phase-packets.html#packet-Protocol::Handshake
|
|
|
|
|
func (mc *mysqlConn) readInitPacket() ([]byte, error) { |
|
|
|
|
func (mc *mysqlConn) readHandshakePacket() ([]byte, string, error) { |
|
|
|
|
data, err := mc.readPacket() |
|
|
|
|
if err != nil { |
|
|
|
|
return nil, err |
|
|
|
|
// for init we can rewrite this to ErrBadConn for sql.Driver to retry, since
|
|
|
|
|
// in connection initialization we don't risk retrying non-idempotent actions.
|
|
|
|
|
if err == ErrInvalidConn { |
|
|
|
|
return nil, "", driver.ErrBadConn |
|
|
|
|
} |
|
|
|
|
return nil, "", err |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
if data[0] == iERR { |
|
|
|
|
return nil, mc.handleErrorPacket(data) |
|
|
|
|
return nil, "", mc.handleErrorPacket(data) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// protocol version [1 byte]
|
|
|
|
|
if data[0] < minProtocolVersion { |
|
|
|
|
return nil, fmt.Errorf( |
|
|
|
|
return nil, "", fmt.Errorf( |
|
|
|
|
"unsupported protocol version %d. Version %d or higher is required", |
|
|
|
|
data[0], |
|
|
|
|
minProtocolVersion, |
|
|
|
@ -157,7 +183,7 @@ func (mc *mysqlConn) readInitPacket() ([]byte, error) { |
|
|
|
|
pos := 1 + bytes.IndexByte(data[1:], 0x00) + 1 + 4 |
|
|
|
|
|
|
|
|
|
// first part of the password cipher [8 bytes]
|
|
|
|
|
cipher := data[pos : pos+8] |
|
|
|
|
authData := data[pos : pos+8] |
|
|
|
|
|
|
|
|
|
// (filler) always 0x00 [1 byte]
|
|
|
|
|
pos += 8 + 1 |
|
|
|
@ -165,13 +191,14 @@ func (mc *mysqlConn) readInitPacket() ([]byte, error) { |
|
|
|
|
// capability flags (lower 2 bytes) [2 bytes]
|
|
|
|
|
mc.flags = clientFlag(binary.LittleEndian.Uint16(data[pos : pos+2])) |
|
|
|
|
if mc.flags&clientProtocol41 == 0 { |
|
|
|
|
return nil, ErrOldProtocol |
|
|
|
|
return nil, "", ErrOldProtocol |
|
|
|
|
} |
|
|
|
|
if mc.flags&clientSSL == 0 && mc.cfg.tls != nil { |
|
|
|
|
return nil, ErrNoTLS |
|
|
|
|
return nil, "", ErrNoTLS |
|
|
|
|
} |
|
|
|
|
pos += 2 |
|
|
|
|
|
|
|
|
|
plugin := "" |
|
|
|
|
if len(data) > pos { |
|
|
|
|
// character set [1 byte]
|
|
|
|
|
// status flags [2 bytes]
|
|
|
|
@ -192,32 +219,34 @@ func (mc *mysqlConn) readInitPacket() ([]byte, error) { |
|
|
|
|
//
|
|
|
|
|
// The official Python library uses the fixed length 12
|
|
|
|
|
// which seems to work but technically could have a hidden bug.
|
|
|
|
|
cipher = append(cipher, data[pos:pos+12]...) |
|
|
|
|
authData = append(authData, data[pos:pos+12]...) |
|
|
|
|
pos += 13 |
|
|
|
|
|
|
|
|
|
// TODO: Verify string termination
|
|
|
|
|
// EOF if version (>= 5.5.7 and < 5.5.10) or (>= 5.6.0 and < 5.6.2)
|
|
|
|
|
// \NUL otherwise
|
|
|
|
|
//
|
|
|
|
|
//if data[len(data)-1] == 0 {
|
|
|
|
|
// return
|
|
|
|
|
//}
|
|
|
|
|
//return ErrMalformPkt
|
|
|
|
|
if end := bytes.IndexByte(data[pos:], 0x00); end != -1 { |
|
|
|
|
plugin = string(data[pos : pos+end]) |
|
|
|
|
} else { |
|
|
|
|
plugin = string(data[pos:]) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// make a memory safe copy of the cipher slice
|
|
|
|
|
var b [20]byte |
|
|
|
|
copy(b[:], cipher) |
|
|
|
|
return b[:], nil |
|
|
|
|
copy(b[:], authData) |
|
|
|
|
return b[:], plugin, nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
plugin = defaultAuthPlugin |
|
|
|
|
|
|
|
|
|
// make a memory safe copy of the cipher slice
|
|
|
|
|
var b [8]byte |
|
|
|
|
copy(b[:], cipher) |
|
|
|
|
return b[:], nil |
|
|
|
|
copy(b[:], authData) |
|
|
|
|
return b[:], plugin, nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Client Authentication Packet
|
|
|
|
|
// http://dev.mysql.com/doc/internals/en/connection-phase-packets.html#packet-Protocol::HandshakeResponse
|
|
|
|
|
func (mc *mysqlConn) writeAuthPacket(cipher []byte) error { |
|
|
|
|
func (mc *mysqlConn) writeHandshakeResponsePacket(authResp []byte, addNUL bool, plugin string) error { |
|
|
|
|
// Adjust client flags based on server support
|
|
|
|
|
clientFlags := clientProtocol41 | |
|
|
|
|
clientSecureConn | |
|
|
|
@ -241,10 +270,19 @@ func (mc *mysqlConn) writeAuthPacket(cipher []byte) error { |
|
|
|
|
clientFlags |= clientMultiStatements |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// User Password
|
|
|
|
|
scrambleBuff := scramblePassword(cipher, []byte(mc.cfg.Passwd)) |
|
|
|
|
// encode length of the auth plugin data
|
|
|
|
|
var authRespLEIBuf [9]byte |
|
|
|
|
authRespLEI := appendLengthEncodedInteger(authRespLEIBuf[:0], uint64(len(authResp))) |
|
|
|
|
if len(authRespLEI) > 1 { |
|
|
|
|
// if the length can not be written in 1 byte, it must be written as a
|
|
|
|
|
// length encoded integer
|
|
|
|
|
clientFlags |= clientPluginAuthLenEncClientData |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
pktLen := 4 + 4 + 1 + 23 + len(mc.cfg.User) + 1 + 1 + len(scrambleBuff) + 21 + 1 |
|
|
|
|
pktLen := 4 + 4 + 1 + 23 + len(mc.cfg.User) + 1 + len(authRespLEI) + len(authResp) + 21 + 1 |
|
|
|
|
if addNUL { |
|
|
|
|
pktLen++ |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// To specify a db name
|
|
|
|
|
if n := len(mc.cfg.DBName); n > 0 { |
|
|
|
@ -257,7 +295,7 @@ func (mc *mysqlConn) writeAuthPacket(cipher []byte) error { |
|
|
|
|
if data == nil { |
|
|
|
|
// cannot take the buffer. Something must be wrong with the connection
|
|
|
|
|
errLog.Print(ErrBusyBuffer) |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
return errBadConnNoWrite |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// ClientFlags [32 bit]
|
|
|
|
@ -312,9 +350,13 @@ func (mc *mysqlConn) writeAuthPacket(cipher []byte) error { |
|
|
|
|
data[pos] = 0x00 |
|
|
|
|
pos++ |
|
|
|
|
|
|
|
|
|
// ScrambleBuffer [length encoded integer]
|
|
|
|
|
data[pos] = byte(len(scrambleBuff)) |
|
|
|
|
pos += 1 + copy(data[pos+1:], scrambleBuff) |
|
|
|
|
// Auth Data [length encoded integer]
|
|
|
|
|
pos += copy(data[pos:], authRespLEI) |
|
|
|
|
pos += copy(data[pos:], authResp) |
|
|
|
|
if addNUL { |
|
|
|
|
data[pos] = 0x00 |
|
|
|
|
pos++ |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Databasename [null terminated string]
|
|
|
|
|
if len(mc.cfg.DBName) > 0 { |
|
|
|
@ -323,72 +365,32 @@ func (mc *mysqlConn) writeAuthPacket(cipher []byte) error { |
|
|
|
|
pos++ |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Assume native client during response
|
|
|
|
|
pos += copy(data[pos:], "mysql_native_password") |
|
|
|
|
pos += copy(data[pos:], plugin) |
|
|
|
|
data[pos] = 0x00 |
|
|
|
|
|
|
|
|
|
// Send Auth packet
|
|
|
|
|
return mc.writePacket(data) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Client old authentication packet
|
|
|
|
|
// http://dev.mysql.com/doc/internals/en/connection-phase-packets.html#packet-Protocol::AuthSwitchResponse
|
|
|
|
|
func (mc *mysqlConn) writeOldAuthPacket(cipher []byte) error { |
|
|
|
|
// User password
|
|
|
|
|
scrambleBuff := scrambleOldPassword(cipher, []byte(mc.cfg.Passwd)) |
|
|
|
|
|
|
|
|
|
// Calculate the packet length and add a tailing 0
|
|
|
|
|
pktLen := len(scrambleBuff) + 1 |
|
|
|
|
data := mc.buf.takeSmallBuffer(4 + pktLen) |
|
|
|
|
if data == nil { |
|
|
|
|
// can not take the buffer. Something must be wrong with the connection
|
|
|
|
|
errLog.Print(ErrBusyBuffer) |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
func (mc *mysqlConn) writeAuthSwitchPacket(authData []byte, addNUL bool) error { |
|
|
|
|
pktLen := 4 + len(authData) |
|
|
|
|
if addNUL { |
|
|
|
|
pktLen++ |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Add the scrambled password [null terminated string]
|
|
|
|
|
copy(data[4:], scrambleBuff) |
|
|
|
|
data[4+pktLen-1] = 0x00 |
|
|
|
|
|
|
|
|
|
return mc.writePacket(data) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Client clear text authentication packet
|
|
|
|
|
// http://dev.mysql.com/doc/internals/en/connection-phase-packets.html#packet-Protocol::AuthSwitchResponse
|
|
|
|
|
func (mc *mysqlConn) writeClearAuthPacket() error { |
|
|
|
|
// Calculate the packet length and add a tailing 0
|
|
|
|
|
pktLen := len(mc.cfg.Passwd) + 1 |
|
|
|
|
data := mc.buf.takeSmallBuffer(4 + pktLen) |
|
|
|
|
data := mc.buf.takeSmallBuffer(pktLen) |
|
|
|
|
if data == nil { |
|
|
|
|
// cannot take the buffer. Something must be wrong with the connection
|
|
|
|
|
errLog.Print(ErrBusyBuffer) |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
return errBadConnNoWrite |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Add the clear password [null terminated string]
|
|
|
|
|
copy(data[4:], mc.cfg.Passwd) |
|
|
|
|
data[4+pktLen-1] = 0x00 |
|
|
|
|
|
|
|
|
|
return mc.writePacket(data) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Native password authentication method
|
|
|
|
|
// http://dev.mysql.com/doc/internals/en/connection-phase-packets.html#packet-Protocol::AuthSwitchResponse
|
|
|
|
|
func (mc *mysqlConn) writeNativeAuthPacket(cipher []byte) error { |
|
|
|
|
scrambleBuff := scramblePassword(cipher, []byte(mc.cfg.Passwd)) |
|
|
|
|
|
|
|
|
|
// Calculate the packet length and add a tailing 0
|
|
|
|
|
pktLen := len(scrambleBuff) |
|
|
|
|
data := mc.buf.takeSmallBuffer(4 + pktLen) |
|
|
|
|
if data == nil { |
|
|
|
|
// can not take the buffer. Something must be wrong with the connection
|
|
|
|
|
errLog.Print(ErrBusyBuffer) |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
// Add the auth data [EOF]
|
|
|
|
|
copy(data[4:], authData) |
|
|
|
|
if addNUL { |
|
|
|
|
data[pktLen-1] = 0x00 |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Add the scramble
|
|
|
|
|
copy(data[4:], scrambleBuff) |
|
|
|
|
|
|
|
|
|
return mc.writePacket(data) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
@ -404,7 +406,7 @@ func (mc *mysqlConn) writeCommandPacket(command byte) error { |
|
|
|
|
if data == nil { |
|
|
|
|
// cannot take the buffer. Something must be wrong with the connection
|
|
|
|
|
errLog.Print(ErrBusyBuffer) |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
return errBadConnNoWrite |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Add command byte
|
|
|
|
@ -423,7 +425,7 @@ func (mc *mysqlConn) writeCommandPacketStr(command byte, arg string) error { |
|
|
|
|
if data == nil { |
|
|
|
|
// cannot take the buffer. Something must be wrong with the connection
|
|
|
|
|
errLog.Print(ErrBusyBuffer) |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
return errBadConnNoWrite |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Add command byte
|
|
|
|
@ -444,7 +446,7 @@ func (mc *mysqlConn) writeCommandPacketUint32(command byte, arg uint32) error { |
|
|
|
|
if data == nil { |
|
|
|
|
// cannot take the buffer. Something must be wrong with the connection
|
|
|
|
|
errLog.Print(ErrBusyBuffer) |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
return errBadConnNoWrite |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Add command byte
|
|
|
|
@ -464,43 +466,50 @@ func (mc *mysqlConn) writeCommandPacketUint32(command byte, arg uint32) error { |
|
|
|
|
* Result Packets * |
|
|
|
|
******************************************************************************/ |
|
|
|
|
|
|
|
|
|
// Returns error if Packet is not an 'Result OK'-Packet
|
|
|
|
|
func (mc *mysqlConn) readResultOK() ([]byte, error) { |
|
|
|
|
func (mc *mysqlConn) readAuthResult() ([]byte, string, error) { |
|
|
|
|
data, err := mc.readPacket() |
|
|
|
|
if err == nil { |
|
|
|
|
if err != nil { |
|
|
|
|
return nil, "", err |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// packet indicator
|
|
|
|
|
switch data[0] { |
|
|
|
|
|
|
|
|
|
case iOK: |
|
|
|
|
return nil, mc.handleOkPacket(data) |
|
|
|
|
return nil, "", mc.handleOkPacket(data) |
|
|
|
|
|
|
|
|
|
case iAuthMoreData: |
|
|
|
|
return data[1:], "", err |
|
|
|
|
|
|
|
|
|
case iEOF: |
|
|
|
|
if len(data) > 1 { |
|
|
|
|
if len(data) < 1 { |
|
|
|
|
// https://dev.mysql.com/doc/internals/en/connection-phase-packets.html#packet-Protocol::OldAuthSwitchRequest
|
|
|
|
|
return nil, "mysql_old_password", nil |
|
|
|
|
} |
|
|
|
|
pluginEndIndex := bytes.IndexByte(data, 0x00) |
|
|
|
|
if pluginEndIndex < 0 { |
|
|
|
|
return nil, "", ErrMalformPkt |
|
|
|
|
} |
|
|
|
|
plugin := string(data[1:pluginEndIndex]) |
|
|
|
|
cipher := data[pluginEndIndex+1 : len(data)-1] |
|
|
|
|
|
|
|
|
|
if plugin == "mysql_old_password" { |
|
|
|
|
// using old_passwords
|
|
|
|
|
return cipher, ErrOldPassword |
|
|
|
|
} else if plugin == "mysql_clear_password" { |
|
|
|
|
// using clear text password
|
|
|
|
|
return cipher, ErrCleartextPassword |
|
|
|
|
} else if plugin == "mysql_native_password" { |
|
|
|
|
// using mysql default authentication method
|
|
|
|
|
return cipher, ErrNativePassword |
|
|
|
|
} else { |
|
|
|
|
return cipher, ErrUnknownPlugin |
|
|
|
|
authData := data[pluginEndIndex+1:] |
|
|
|
|
return authData, plugin, nil |
|
|
|
|
|
|
|
|
|
default: // Error otherwise
|
|
|
|
|
return nil, "", mc.handleErrorPacket(data) |
|
|
|
|
} |
|
|
|
|
} else { |
|
|
|
|
return nil, ErrOldPassword |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
default: // Error otherwise
|
|
|
|
|
return nil, mc.handleErrorPacket(data) |
|
|
|
|
// Returns error if Packet is not an 'Result OK'-Packet
|
|
|
|
|
func (mc *mysqlConn) readResultOK() error { |
|
|
|
|
data, err := mc.readPacket() |
|
|
|
|
if err != nil { |
|
|
|
|
return err |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
if data[0] == iOK { |
|
|
|
|
return mc.handleOkPacket(data) |
|
|
|
|
} |
|
|
|
|
return nil, err |
|
|
|
|
return mc.handleErrorPacket(data) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Result Set Header Packet
|
|
|
|
@ -543,6 +552,22 @@ func (mc *mysqlConn) handleErrorPacket(data []byte) error { |
|
|
|
|
// Error Number [16 bit uint]
|
|
|
|
|
errno := binary.LittleEndian.Uint16(data[1:3]) |
|
|
|
|
|
|
|
|
|
// 1792: ER_CANT_EXECUTE_IN_READ_ONLY_TRANSACTION
|
|
|
|
|
// 1290: ER_OPTION_PREVENTS_STATEMENT (returned by Aurora during failover)
|
|
|
|
|
if (errno == 1792 || errno == 1290) && mc.cfg.RejectReadOnly { |
|
|
|
|
// Oops; we are connected to a read-only connection, and won't be able
|
|
|
|
|
// to issue any write statements. Since RejectReadOnly is configured,
|
|
|
|
|
// we throw away this connection hoping this one would have write
|
|
|
|
|
// permission. This is specifically for a possible race condition
|
|
|
|
|
// during failover (e.g. on AWS Aurora). See README.md for more.
|
|
|
|
|
//
|
|
|
|
|
// We explicitly close the connection before returning
|
|
|
|
|
// driver.ErrBadConn to ensure that `database/sql` purges this
|
|
|
|
|
// connection and initiates a new one for next statement next time.
|
|
|
|
|
mc.Close() |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
pos := 3 |
|
|
|
|
|
|
|
|
|
// SQL State [optional: # + 5bytes string]
|
|
|
|
@ -577,19 +602,12 @@ func (mc *mysqlConn) handleOkPacket(data []byte) error { |
|
|
|
|
|
|
|
|
|
// server_status [2 bytes]
|
|
|
|
|
mc.status = readStatus(data[1+n+m : 1+n+m+2]) |
|
|
|
|
if err := mc.discardResults(); err != nil { |
|
|
|
|
return err |
|
|
|
|
if mc.status&statusMoreResultsExists != 0 { |
|
|
|
|
return nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// warning count [2 bytes]
|
|
|
|
|
if !mc.strict { |
|
|
|
|
return nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
pos := 1 + n + m + 2 |
|
|
|
|
if binary.LittleEndian.Uint16(data[pos:pos+2]) > 0 { |
|
|
|
|
return mc.getWarnings() |
|
|
|
|
} |
|
|
|
|
return nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
@ -661,14 +679,21 @@ func (mc *mysqlConn) readColumns(count int) ([]mysqlField, error) { |
|
|
|
|
if err != nil { |
|
|
|
|
return nil, err |
|
|
|
|
} |
|
|
|
|
pos += n |
|
|
|
|
|
|
|
|
|
// Filler [uint8]
|
|
|
|
|
pos++ |
|
|
|
|
|
|
|
|
|
// Charset [charset, collation uint8]
|
|
|
|
|
columns[i].charSet = data[pos] |
|
|
|
|
pos += 2 |
|
|
|
|
|
|
|
|
|
// Length [uint32]
|
|
|
|
|
pos += n + 1 + 2 + 4 |
|
|
|
|
columns[i].length = binary.LittleEndian.Uint32(data[pos : pos+4]) |
|
|
|
|
pos += 4 |
|
|
|
|
|
|
|
|
|
// Field type [uint8]
|
|
|
|
|
columns[i].fieldType = data[pos] |
|
|
|
|
columns[i].fieldType = fieldType(data[pos]) |
|
|
|
|
pos++ |
|
|
|
|
|
|
|
|
|
// Flags [uint16]
|
|
|
|
@ -691,6 +716,10 @@ func (mc *mysqlConn) readColumns(count int) ([]mysqlField, error) { |
|
|
|
|
func (rows *textRows) readRow(dest []driver.Value) error { |
|
|
|
|
mc := rows.mc |
|
|
|
|
|
|
|
|
|
if rows.rs.done { |
|
|
|
|
return io.EOF |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
data, err := mc.readPacket() |
|
|
|
|
if err != nil { |
|
|
|
|
return err |
|
|
|
@ -700,10 +729,10 @@ func (rows *textRows) readRow(dest []driver.Value) error { |
|
|
|
|
if data[0] == iEOF && len(data) == 5 { |
|
|
|
|
// server_status [2 bytes]
|
|
|
|
|
rows.mc.status = readStatus(data[3:]) |
|
|
|
|
if err := rows.mc.discardResults(); err != nil { |
|
|
|
|
return err |
|
|
|
|
} |
|
|
|
|
rows.rs.done = true |
|
|
|
|
if !rows.HasNextResultSet() { |
|
|
|
|
rows.mc = nil |
|
|
|
|
} |
|
|
|
|
return io.EOF |
|
|
|
|
} |
|
|
|
|
if data[0] == iERR { |
|
|
|
@ -725,7 +754,7 @@ func (rows *textRows) readRow(dest []driver.Value) error { |
|
|
|
|
if !mc.parseTime { |
|
|
|
|
continue |
|
|
|
|
} else { |
|
|
|
|
switch rows.columns[i].fieldType { |
|
|
|
|
switch rows.rs.columns[i].fieldType { |
|
|
|
|
case fieldTypeTimestamp, fieldTypeDateTime, |
|
|
|
|
fieldTypeDate, fieldTypeNewDate: |
|
|
|
|
dest[i], err = parseDateTime( |
|
|
|
@ -797,14 +826,7 @@ func (stmt *mysqlStmt) readPrepareResultPacket() (uint16, error) { |
|
|
|
|
// Reserved [8 bit]
|
|
|
|
|
|
|
|
|
|
// Warning count [16 bit uint]
|
|
|
|
|
if !stmt.mc.strict { |
|
|
|
|
return columnCount, nil |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Check for warnings count > 0, only available in MySQL > 4.1
|
|
|
|
|
if len(data) >= 12 && binary.LittleEndian.Uint16(data[10:12]) > 0 { |
|
|
|
|
return columnCount, stmt.mc.getWarnings() |
|
|
|
|
} |
|
|
|
|
return columnCount, nil |
|
|
|
|
} |
|
|
|
|
return 0, err |
|
|
|
@ -876,6 +898,12 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
const minPktLen = 4 + 1 + 4 + 1 + 4 |
|
|
|
|
mc := stmt.mc |
|
|
|
|
|
|
|
|
|
// Determine threshould dynamically to avoid packet size shortage.
|
|
|
|
|
longDataSize := mc.maxAllowedPacket / (stmt.paramCount + 1) |
|
|
|
|
if longDataSize < 64 { |
|
|
|
|
longDataSize = 64 |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Reset packet-sequence
|
|
|
|
|
mc.sequence = 0 |
|
|
|
|
|
|
|
|
@ -889,7 +917,7 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
if data == nil { |
|
|
|
|
// cannot take the buffer. Something must be wrong with the connection
|
|
|
|
|
errLog.Print(ErrBusyBuffer) |
|
|
|
|
return driver.ErrBadConn |
|
|
|
|
return errBadConnNoWrite |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// command [1 byte]
|
|
|
|
@ -948,7 +976,7 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
// build NULL-bitmap
|
|
|
|
|
if arg == nil { |
|
|
|
|
nullMask[i/8] |= 1 << (uint(i) & 7) |
|
|
|
|
paramTypes[i+i] = fieldTypeNULL |
|
|
|
|
paramTypes[i+i] = byte(fieldTypeNULL) |
|
|
|
|
paramTypes[i+i+1] = 0x00 |
|
|
|
|
continue |
|
|
|
|
} |
|
|
|
@ -956,7 +984,7 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
// cache types and values
|
|
|
|
|
switch v := arg.(type) { |
|
|
|
|
case int64: |
|
|
|
|
paramTypes[i+i] = fieldTypeLongLong |
|
|
|
|
paramTypes[i+i] = byte(fieldTypeLongLong) |
|
|
|
|
paramTypes[i+i+1] = 0x00 |
|
|
|
|
|
|
|
|
|
if cap(paramValues)-len(paramValues)-8 >= 0 { |
|
|
|
@ -972,7 +1000,7 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
case float64: |
|
|
|
|
paramTypes[i+i] = fieldTypeDouble |
|
|
|
|
paramTypes[i+i] = byte(fieldTypeDouble) |
|
|
|
|
paramTypes[i+i+1] = 0x00 |
|
|
|
|
|
|
|
|
|
if cap(paramValues)-len(paramValues)-8 >= 0 { |
|
|
|
@ -988,7 +1016,7 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
case bool: |
|
|
|
|
paramTypes[i+i] = fieldTypeTiny |
|
|
|
|
paramTypes[i+i] = byte(fieldTypeTiny) |
|
|
|
|
paramTypes[i+i+1] = 0x00 |
|
|
|
|
|
|
|
|
|
if v { |
|
|
|
@ -1000,10 +1028,10 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
case []byte: |
|
|
|
|
// Common case (non-nil value) first
|
|
|
|
|
if v != nil { |
|
|
|
|
paramTypes[i+i] = fieldTypeString |
|
|
|
|
paramTypes[i+i] = byte(fieldTypeString) |
|
|
|
|
paramTypes[i+i+1] = 0x00 |
|
|
|
|
|
|
|
|
|
if len(v) < mc.maxAllowedPacket-pos-len(paramValues)-(len(args)-(i+1))*64 { |
|
|
|
|
if len(v) < longDataSize { |
|
|
|
|
paramValues = appendLengthEncodedInteger(paramValues, |
|
|
|
|
uint64(len(v)), |
|
|
|
|
) |
|
|
|
@ -1018,14 +1046,14 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
|
|
|
|
|
// Handle []byte(nil) as a NULL value
|
|
|
|
|
nullMask[i/8] |= 1 << (uint(i) & 7) |
|
|
|
|
paramTypes[i+i] = fieldTypeNULL |
|
|
|
|
paramTypes[i+i] = byte(fieldTypeNULL) |
|
|
|
|
paramTypes[i+i+1] = 0x00 |
|
|
|
|
|
|
|
|
|
case string: |
|
|
|
|
paramTypes[i+i] = fieldTypeString |
|
|
|
|
paramTypes[i+i] = byte(fieldTypeString) |
|
|
|
|
paramTypes[i+i+1] = 0x00 |
|
|
|
|
|
|
|
|
|
if len(v) < mc.maxAllowedPacket-pos-len(paramValues)-(len(args)-(i+1))*64 { |
|
|
|
|
if len(v) < longDataSize { |
|
|
|
|
paramValues = appendLengthEncodedInteger(paramValues, |
|
|
|
|
uint64(len(v)), |
|
|
|
|
) |
|
|
|
@ -1037,20 +1065,22 @@ func (stmt *mysqlStmt) writeExecutePacket(args []driver.Value) error { |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
case time.Time: |
|
|
|
|
paramTypes[i+i] = fieldTypeString |
|
|
|
|
paramTypes[i+i] = byte(fieldTypeString) |
|
|
|
|
paramTypes[i+i+1] = 0x00 |
|
|
|
|
|
|
|
|
|
var val []byte |
|
|
|
|
var a [64]byte |
|
|
|
|
var b = a[:0] |
|
|
|
|
|
|
|
|
|
if v.IsZero() { |
|
|
|
|
val = []byte("0000-00-00") |
|
|
|
|
b = append(b, "0000-00-00"...) |
|
|
|
|
} else { |
|
|
|
|
val = []byte(v.In(mc.cfg.Loc).Format(timeFormat)) |
|
|
|
|
b = v.In(mc.cfg.Loc).AppendFormat(b, timeFormat) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
paramValues = appendLengthEncodedInteger(paramValues, |
|
|
|
|
uint64(len(val)), |
|
|
|
|
uint64(len(b)), |
|
|
|
|
) |
|
|
|
|
paramValues = append(paramValues, val...) |
|
|
|
|
paramValues = append(paramValues, b...) |
|
|
|
|
|
|
|
|
|
default: |
|
|
|
|
return fmt.Errorf("cannot convert type: %T", arg) |
|
|
|
@ -1086,8 +1116,6 @@ func (mc *mysqlConn) discardResults() error { |
|
|
|
|
if err := mc.readUntilEOF(); err != nil { |
|
|
|
|
return err |
|
|
|
|
} |
|
|
|
|
} else { |
|
|
|
|
mc.status &^= statusMoreResultsExists |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
return nil |
|
|
|
@ -1105,16 +1133,17 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
// EOF Packet
|
|
|
|
|
if data[0] == iEOF && len(data) == 5 { |
|
|
|
|
rows.mc.status = readStatus(data[3:]) |
|
|
|
|
if err := rows.mc.discardResults(); err != nil { |
|
|
|
|
return err |
|
|
|
|
} |
|
|
|
|
rows.rs.done = true |
|
|
|
|
if !rows.HasNextResultSet() { |
|
|
|
|
rows.mc = nil |
|
|
|
|
} |
|
|
|
|
return io.EOF |
|
|
|
|
} |
|
|
|
|
mc := rows.mc |
|
|
|
|
rows.mc = nil |
|
|
|
|
|
|
|
|
|
// Error otherwise
|
|
|
|
|
return rows.mc.handleErrorPacket(data) |
|
|
|
|
return mc.handleErrorPacket(data) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// NULL-bitmap, [(column-count + 7 + 2) / 8 bytes]
|
|
|
|
@ -1130,14 +1159,14 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Convert to byte-coded string
|
|
|
|
|
switch rows.columns[i].fieldType { |
|
|
|
|
switch rows.rs.columns[i].fieldType { |
|
|
|
|
case fieldTypeNULL: |
|
|
|
|
dest[i] = nil |
|
|
|
|
continue |
|
|
|
|
|
|
|
|
|
// Numeric Types
|
|
|
|
|
case fieldTypeTiny: |
|
|
|
|
if rows.columns[i].flags&flagUnsigned != 0 { |
|
|
|
|
if rows.rs.columns[i].flags&flagUnsigned != 0 { |
|
|
|
|
dest[i] = int64(data[pos]) |
|
|
|
|
} else { |
|
|
|
|
dest[i] = int64(int8(data[pos])) |
|
|
|
@ -1146,7 +1175,7 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
continue |
|
|
|
|
|
|
|
|
|
case fieldTypeShort, fieldTypeYear: |
|
|
|
|
if rows.columns[i].flags&flagUnsigned != 0 { |
|
|
|
|
if rows.rs.columns[i].flags&flagUnsigned != 0 { |
|
|
|
|
dest[i] = int64(binary.LittleEndian.Uint16(data[pos : pos+2])) |
|
|
|
|
} else { |
|
|
|
|
dest[i] = int64(int16(binary.LittleEndian.Uint16(data[pos : pos+2]))) |
|
|
|
@ -1155,7 +1184,7 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
continue |
|
|
|
|
|
|
|
|
|
case fieldTypeInt24, fieldTypeLong: |
|
|
|
|
if rows.columns[i].flags&flagUnsigned != 0 { |
|
|
|
|
if rows.rs.columns[i].flags&flagUnsigned != 0 { |
|
|
|
|
dest[i] = int64(binary.LittleEndian.Uint32(data[pos : pos+4])) |
|
|
|
|
} else { |
|
|
|
|
dest[i] = int64(int32(binary.LittleEndian.Uint32(data[pos : pos+4]))) |
|
|
|
@ -1164,7 +1193,7 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
continue |
|
|
|
|
|
|
|
|
|
case fieldTypeLongLong: |
|
|
|
|
if rows.columns[i].flags&flagUnsigned != 0 { |
|
|
|
|
if rows.rs.columns[i].flags&flagUnsigned != 0 { |
|
|
|
|
val := binary.LittleEndian.Uint64(data[pos : pos+8]) |
|
|
|
|
if val > math.MaxInt64 { |
|
|
|
|
dest[i] = uint64ToString(val) |
|
|
|
@ -1178,7 +1207,7 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
continue |
|
|
|
|
|
|
|
|
|
case fieldTypeFloat: |
|
|
|
|
dest[i] = float32(math.Float32frombits(binary.LittleEndian.Uint32(data[pos : pos+4]))) |
|
|
|
|
dest[i] = math.Float32frombits(binary.LittleEndian.Uint32(data[pos : pos+4])) |
|
|
|
|
pos += 4 |
|
|
|
|
continue |
|
|
|
|
|
|
|
|
@ -1218,10 +1247,10 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
case isNull: |
|
|
|
|
dest[i] = nil |
|
|
|
|
continue |
|
|
|
|
case rows.columns[i].fieldType == fieldTypeTime: |
|
|
|
|
case rows.rs.columns[i].fieldType == fieldTypeTime: |
|
|
|
|
// database/sql does not support an equivalent to TIME, return a string
|
|
|
|
|
var dstlen uint8 |
|
|
|
|
switch decimals := rows.columns[i].decimals; decimals { |
|
|
|
|
switch decimals := rows.rs.columns[i].decimals; decimals { |
|
|
|
|
case 0x00, 0x1f: |
|
|
|
|
dstlen = 8 |
|
|
|
|
case 1, 2, 3, 4, 5, 6: |
|
|
|
@ -1229,7 +1258,7 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
default: |
|
|
|
|
return fmt.Errorf( |
|
|
|
|
"protocol error, illegal decimals value %d", |
|
|
|
|
rows.columns[i].decimals, |
|
|
|
|
rows.rs.columns[i].decimals, |
|
|
|
|
) |
|
|
|
|
} |
|
|
|
|
dest[i], err = formatBinaryDateTime(data[pos:pos+int(num)], dstlen, true) |
|
|
|
@ -1237,10 +1266,10 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
dest[i], err = parseBinaryDateTime(num, data[pos:], rows.mc.cfg.Loc) |
|
|
|
|
default: |
|
|
|
|
var dstlen uint8 |
|
|
|
|
if rows.columns[i].fieldType == fieldTypeDate { |
|
|
|
|
if rows.rs.columns[i].fieldType == fieldTypeDate { |
|
|
|
|
dstlen = 10 |
|
|
|
|
} else { |
|
|
|
|
switch decimals := rows.columns[i].decimals; decimals { |
|
|
|
|
switch decimals := rows.rs.columns[i].decimals; decimals { |
|
|
|
|
case 0x00, 0x1f: |
|
|
|
|
dstlen = 19 |
|
|
|
|
case 1, 2, 3, 4, 5, 6: |
|
|
|
@ -1248,7 +1277,7 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
default: |
|
|
|
|
return fmt.Errorf( |
|
|
|
|
"protocol error, illegal decimals value %d", |
|
|
|
|
rows.columns[i].decimals, |
|
|
|
|
rows.rs.columns[i].decimals, |
|
|
|
|
) |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
@ -1264,7 +1293,7 @@ func (rows *binaryRows) readRow(dest []driver.Value) error { |
|
|
|
|
|
|
|
|
|
// Please report if this happens!
|
|
|
|
|
default: |
|
|
|
|
return fmt.Errorf("unknown field type %d", rows.columns[i].fieldType) |
|
|
|
|
return fmt.Errorf("unknown field type %d", rows.rs.columns[i].fieldType) |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|