| 2496 | } |
| 2497 | |
| 2498 | bool readSnapshotBulkPayload(connection *conn, redisMaster *mi, rdbSaveInfo &rsi) { |
| 2499 | int fUpdate = g_pserver->fActiveReplica || g_pserver->enable_multimaster; |
| 2500 | serverAssert(GlobalLocksAcquired()); |
| 2501 | serverAssert(mi->master == nullptr); |
| 2502 | bool fFinished = false; |
| 2503 | |
| 2504 | if (mi->bulkreadBuffer == nullptr) { |
| 2505 | mi->bulkreadBuffer = sdsempty(); |
| 2506 | mi->parseState = new SnapshotPayloadParseState(); |
| 2507 | if (g_pserver->aof_state != AOF_OFF) stopAppendOnly(); |
| 2508 | if (!fUpdate) { |
| 2509 | int empty_db_flags = g_pserver->repl_slave_lazy_flush ? EMPTYDB_ASYNC : |
| 2510 | EMPTYDB_NO_FLAGS; |
| 2511 | serverLog(LL_NOTICE, "MASTER <-> REPLICA sync: Flushing old data"); |
| 2512 | emptyDb(-1,empty_db_flags,replicationEmptyDbCallback); |
| 2513 | for (int idb = 0; idb < cserver.dbnum; ++idb) { |
| 2514 | aeAcquireLock(); |
| 2515 | g_pserver->db[idb]->processChanges(false); |
| 2516 | aeReleaseLock(); |
| 2517 | g_pserver->db[idb]->commitChanges(); |
| 2518 | g_pserver->db[idb]->trackChanges(false); |
| 2519 | } |
| 2520 | } |
| 2521 | } |
| 2522 | |
| 2523 | serverAssert(mi->parseState != nullptr); |
| 2524 | for (int iter = 0; iter < 10; ++iter) { |
| 2525 | if (mi->parseState->shouldThrottle()) |
| 2526 | return false; |
| 2527 | |
| 2528 | auto readlen = PROTO_IOBUF_LEN; |
| 2529 | auto qblen = sdslen(mi->bulkreadBuffer); |
| 2530 | mi->bulkreadBuffer = sdsMakeRoomFor(mi->bulkreadBuffer, readlen); |
| 2531 | |
| 2532 | auto nread = connRead(conn, mi->bulkreadBuffer+qblen, readlen); |
| 2533 | if (nread <= 0) { |
| 2534 | if (connGetState(conn) != CONN_STATE_CONNECTED) { |
| 2535 | serverLog(LL_WARNING,"I/O error trying to sync with MASTER: %s", |
| 2536 | (nread == -1) ? strerror(errno) : "connection lost"); |
| 2537 | cancelReplicationHandshake(mi, true); |
| 2538 | } |
| 2539 | return false; |
| 2540 | } |
| 2541 | g_pserver->stat_net_input_bytes += nread; |
| 2542 | mi->repl_transfer_lastio = g_pserver->unixtime; |
| 2543 | mi->repl_transfer_read += nread; |
| 2544 | sdsIncrLen(mi->bulkreadBuffer,nread); |
| 2545 | |
| 2546 | size_t offset = 0; |
| 2547 | |
| 2548 | try { |
| 2549 | if (sdslen(mi->bulkreadBuffer) > cserver.client_max_querybuf_len) { |
| 2550 | throw "Full Sync Streaming Buffer Exceeded (increase client_max_querybuf_len)"; |
| 2551 | } |
| 2552 | |
| 2553 | while (sdslen(mi->bulkreadBuffer) > offset) { |
| 2554 | // Pop completed items |
| 2555 | mi->parseState->trimState(); |
no test coverage detected