diff --git a/index.js b/index.js index db6746f9..dd1d8593 100644 --- a/index.js +++ b/index.js @@ -157,12 +157,14 @@ ZongJi.prototype._options = function({ filename, position, startAtEnd, + gtidsData, }) { this.options = { serverId, filename, position, startAtEnd, + gtidsData, }; }; @@ -204,6 +206,7 @@ ZongJi.prototype.get = function(name) { // - `filename`, `position` the position of binlog to beigin with // - `startAtEnd` if true, will update filename / postion automatically // - `includeEvents`, `excludeEvents`, `includeSchema`, `exludeSchema` filter different binlog events bubbling +// - `gtidsData` {[SID]: Array<[start, end]>} GTIDs-set with processed events GTIDs for server-side events filtering ZongJi.prototype.start = function(options = {}) { this._options(options); diff --git a/lib/binlog_event.js b/lib/binlog_event.js index 1a905219..25ca323b 100644 --- a/lib/binlog_event.js +++ b/lib/binlog_event.js @@ -89,7 +89,7 @@ function Query(parser) { this.errorCode = parser.parseUnsignedNumber(2); this.statusVarsLength = parser.parseUnsignedNumber(2); - this.statusVars = parser.parseString(this.statusVarsLength); + this.statusVars = parser.parseBuffer(this.statusVarsLength); this.schema = parser.parseString(this.schemaLength); parser.parseUnsignedNumber(1); @@ -285,11 +285,40 @@ TableMap.prototype.dump = function() { console.log('Column types:', this.columnTypes); }; +/** + * GTID_LOG_EVENT + * see https://dev.mysql.com/doc/dev/mysql-server/latest/classbinary__log_1_1Gtid__event.html + **/ +function GTIDLog(parser) { + BinlogEvent.apply(this, arguments); + this.GTID_FLAGS = parser.parseUnsignedNumber(1); + + const SIDBuffer = parser._buffer.slice(parser._offset, parser._offset + 16); + parser._offset += 16; + this.SID = SIDBuffer.toString('hex'); + + this.GNO = BigInt(parser.parseUnsignedNumber(4)); + this.GNO += BigInt(parser.parseUnsignedNumber(4)) << BigInt(32); +} +util.inherits(GTIDLog, BinlogEvent); + +GTIDLog.prototype.dump = function() { + BinlogEvent.prototype.dump.apply(this); + console.log('SID: %s', this.SID); + console.log('GNO: %d', this.GNO); +}; + + function Unknown() { BinlogEvent.apply(this, arguments); } util.inherits(Unknown, BinlogEvent); +function HeartBeat() { + BinlogEvent.apply(this, arguments); +} +util.inherits(HeartBeat, BinlogEvent); + exports.BinlogEvent = BinlogEvent; exports.Rotate = Rotate; exports.Format = Format; @@ -297,4 +326,6 @@ exports.Query = Query; exports.IntVar = IntVar; exports.Xid = Xid; exports.TableMap = TableMap; +exports.GTIDLog = GTIDLog; +exports.HeartBeat = HeartBeat; exports.Unknown = Unknown; diff --git a/lib/code_map.js b/lib/code_map.js index 1e66b884..871f8b7c 100644 --- a/lib/code_map.js +++ b/lib/code_map.js @@ -47,6 +47,8 @@ const EventClass = { ROTATE_EVENT: events.Rotate, FORMAT_DESCRIPTION_EVENT: events.Format, XID_EVENT: events.Xid, + GTID_LOG_EVENT: events.GTIDLog, + HEARTBEAT_LOG_EVENT: events.HeartBeat, TABLE_MAP_EVENT: events.TableMap, DELETE_ROWS_EVENT_V1: rowsEvents.DeleteRows, diff --git a/lib/packet/combinloggtid.js b/lib/packet/combinloggtid.js new file mode 100644 index 00000000..9dca4c08 --- /dev/null +++ b/lib/packet/combinloggtid.js @@ -0,0 +1,75 @@ +const invariant = require('invariant'); + +function ComBinlogGTID({serverId, nonBlock, filename, position, gtidsData}) { + this.command = 0x1e; + this.position = position || 4; + + this.flags = 0; + this.flags |= (nonBlock ? 1 : 0); + this.flags |= 0x04; /*BINLOG_THROUGH_GTID*/ + + this.serverId = serverId || 1; + this.filename = filename || ''; + + this.gtidsData = gtidsData || {}; +} + +const writeUInt64 = (writer, _value /* BigInt */) => { + const value = BigInt(_value); + + writer._allocate(8); + + for (let i = 0; i < 8; i++) { + writer._buffer[writer._offset++] = Number((value >> BigInt(i * 8)) & BigInt(0xff)); + } +}; + +/** + * https://dev.mysql.com/doc/internals/en/com-binlog-dump-gtid.html + */ +ComBinlogGTID.prototype.write = function(writer) { + writer.writeUnsignedNumber(1, this.command); + writer.writeUnsignedNumber(2, this.flags); + writer.writeUnsignedNumber(4, this.serverId); + + writer.writeUnsignedNumber(4, Buffer.byteLength(this.filename, 'utf-8')); + writer.writeString(this.filename); + + writer.writeUnsignedNumber(4, this.position); + writer.writeUnsignedNumber(4, 0); // high-part of this.position + + if (!(this.flags & 0x04)) return; + + const gtidsDataEntries = Object.entries(this.gtidsData); + // TODO: Support for multiple intervals per SID + writer.writeUnsignedNumber(4, 8 + gtidsDataEntries.length * 40); //data_length + + writer.writeUnsignedNumber(4, gtidsDataEntries.length); //n_sids + writer.writeUnsignedNumber(4, 0); //n_sids + + for (let i = 0; i < gtidsDataEntries.length; i++) { + const [sid, intervals] = gtidsDataEntries[i]; + + const sidBuffer = Buffer.from(sid, 'hex'); + invariant(sidBuffer.length === 16, 'SID should be 16-bytes hex-string'); + writer.writeBuffer(sidBuffer); // sid + // Buffer.byteLength(value, 'utf-8') + + writer.writeUnsignedNumber(4, intervals.length); // n_intervals + writer.writeUnsignedNumber(4, 0); // high-part of n_intervals + + for (let j = 0; j < intervals.length; j++) { + const [start, end] = intervals[j]; + writeUInt64(writer, start); + writeUInt64(writer, end); + } + } + console.log("!!!ComBinlogGTID::write", this, writer._buffer, writer._offset); + console.log(writer._buffer.toString('hex')); +}; + +ComBinlogGTID.prototype.parse = function() { + throw new Error('should never be called here'); +}; + +module.exports = ComBinlogGTID; diff --git a/lib/packet/index.js b/lib/packet/index.js index 3f9f7da4..c07fadcf 100644 --- a/lib/packet/index.js +++ b/lib/packet/index.js @@ -63,4 +63,5 @@ ErrorPacket.prototype.write = function(writer) { exports.EofPacket = EofPacket; exports.ErrorPacket = ErrorPacket; exports.ComBinlog = require('./combinlog'); +exports.ComBinlogGTID = require('./combinloggtid'); exports.initBinlogPacketClass = require('./binlog'); diff --git a/lib/sequence/binlog.js b/lib/sequence/binlog.js index 74bff202..a828a9f0 100644 --- a/lib/sequence/binlog.js +++ b/lib/sequence/binlog.js @@ -1,5 +1,5 @@ const Util = require('util'); -const { EofPacket, ErrorPacket, ComBinlog, initBinlogPacketClass } = require('../packet'); +const { EofPacket, ErrorPacket, ComBinlog, ComBinlogGTID, initBinlogPacketClass } = require('../packet'); const Sequence = require('mysql/lib/protocol/sequences').Sequence; module.exports = function(zongji) { @@ -14,9 +14,12 @@ module.exports = function(zongji) { Binlog.prototype.start = function() { // options include: position / nonBlock / serverId / filename let options = zongji.get([ - 'serverId', 'position', 'filename', 'nonBlock', + 'serverId', 'position', 'filename', 'nonBlock', 'gtidsData' ]); - this.emit('packet', new ComBinlog(options)); + + const ComBinlogClass = options.gtidsData ? ComBinlogGTID : ComBinlog; + + this.emit('packet', new ComBinlogClass(options)); }; Binlog.prototype.determinePacket = function(firstByte) { diff --git a/package-lock.json b/package-lock.json index 8ecd07a5..d09441c0 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1172,6 +1172,14 @@ "through": "^2.3.6" } }, + "invariant": { + "version": "2.2.4", + "resolved": "https://registry.npmjs.org/invariant/-/invariant-2.2.4.tgz", + "integrity": "sha512-phJfQVBuaJM5raOpJjSfkiD6BpbCE4Ns//LaXl6wGYtUBY83nWS6Rf9tXm2e8VaK60JEjYldbPif/A2B1C2gNA==", + "requires": { + "loose-envify": "^1.0.0" + } + }, "is-arrayish": { "version": "0.2.1", "resolved": "https://registry.npmjs.org/is-arrayish/-/is-arrayish-0.2.1.tgz", @@ -1348,8 +1356,7 @@ "js-tokens": { "version": "4.0.0", "resolved": "https://registry.npmjs.org/js-tokens/-/js-tokens-4.0.0.tgz", - "integrity": "sha512-RdJUflcE3cUzKiMqQgsCu06FPu9UdIJO0beYbPhHN4k6apgJtifcoCtT9bcxOpYBtpD2kCM6Sbzg4CausW/PKQ==", - "dev": true + "integrity": "sha512-RdJUflcE3cUzKiMqQgsCu06FPu9UdIJO0beYbPhHN4k6apgJtifcoCtT9bcxOpYBtpD2kCM6Sbzg4CausW/PKQ==" }, "js-yaml": { "version": "3.13.1", @@ -1479,6 +1486,14 @@ "integrity": "sha512-U7KCmLdqsGHBLeWqYlFA0V0Sl6P08EE1ZrmA9cxjUE0WVqT9qnyVDPz1kzpFEP0jdJuFnasWIfSd7fsaNXkpbg==", "dev": true }, + "loose-envify": { + "version": "1.4.0", + "resolved": "https://registry.npmjs.org/loose-envify/-/loose-envify-1.4.0.tgz", + "integrity": "sha512-lyuxPGr/Wfhrlem2CL/UcnUc1zcqKAImBDzukY7Y5F/yQiNdko6+fRLevlw1HgMySw7f611UIY408EtxRSoK3Q==", + "requires": { + "js-tokens": "^3.0.0 || ^4.0.0" + } + }, "lru-cache": { "version": "4.1.5", "resolved": "https://registry.npmjs.org/lru-cache/-/lru-cache-4.1.5.tgz", diff --git a/package.json b/package.json index 2c6f391b..e93597ae 100644 --- a/package.json +++ b/package.json @@ -38,6 +38,7 @@ "dependencies": { "big-integer": "1.6.47", "iconv-lite": "0.4.24", + "invariant": "^2.2.2", "mysql": "2.17.1" } }