| 25 | #include "flow/actorcompiler.h" // This must be the last #include. |
| 26 | |
| 27 | class DDTxnProcessorImpl { |
| 28 | friend class DDTxnProcessor; |
| 29 | |
| 30 | // return {sourceServers, completeSources} |
| 31 | ACTOR static Future<IDDTxnProcessor::SourceServers> getSourceServersForRange(Database cx, KeyRangeRef keys) { |
| 32 | state std::set<UID> servers; |
| 33 | state std::vector<UID> completeSources; |
| 34 | state Transaction tr(cx); |
| 35 | |
| 36 | loop { |
| 37 | servers.clear(); |
| 38 | completeSources.clear(); |
| 39 | |
| 40 | tr.setOption(FDBTransactionOptions::READ_SYSTEM_KEYS); |
| 41 | tr.setOption(FDBTransactionOptions::READ_LOCK_AWARE); |
| 42 | tr.setOption(FDBTransactionOptions::PRIORITY_SYSTEM_IMMEDIATE); |
| 43 | try { |
| 44 | state RangeResult UIDtoTagMap = wait(tr.getRange(serverTagKeys, CLIENT_KNOBS->TOO_MANY)); |
| 45 | ASSERT(!UIDtoTagMap.more && UIDtoTagMap.size() < CLIENT_KNOBS->TOO_MANY); |
| 46 | RangeResult keyServersEntries = wait(tr.getRange(lastLessOrEqual(keyServersKey(keys.begin)), |
| 47 | firstGreaterOrEqual(keyServersKey(keys.end)), |
| 48 | SERVER_KNOBS->DD_QUEUE_MAX_KEY_SERVERS)); |
| 49 | |
| 50 | if (keyServersEntries.size() < SERVER_KNOBS->DD_QUEUE_MAX_KEY_SERVERS) { |
| 51 | for (int shard = 0; shard < keyServersEntries.size(); shard++) { |
| 52 | std::vector<UID> src, dest; |
| 53 | decodeKeyServersValue(UIDtoTagMap, keyServersEntries[shard].value, src, dest); |
| 54 | ASSERT(src.size()); |
| 55 | for (int i = 0; i < src.size(); i++) { |
| 56 | servers.insert(src[i]); |
| 57 | } |
| 58 | if (shard == 0) { |
| 59 | completeSources = src; |
| 60 | } else { |
| 61 | for (int i = 0; i < completeSources.size(); i++) { |
| 62 | if (std::find(src.begin(), src.end(), completeSources[i]) == src.end()) { |
| 63 | swapAndPop(&completeSources, i--); |
| 64 | } |
| 65 | } |
| 66 | } |
| 67 | } |
| 68 | |
| 69 | ASSERT(servers.size() > 0); |
| 70 | } |
| 71 | |
| 72 | // If the size of keyServerEntries is large, then just assume we are using all storage servers |
| 73 | // Why the size can be large? |
| 74 | // When a shard is inflight and DD crashes, some destination servers may have already got the data. |
| 75 | // The new DD will treat the destination servers as source servers. So the size can be large. |
| 76 | else { |
| 77 | RangeResult serverList = wait(tr.getRange(serverListKeys, CLIENT_KNOBS->TOO_MANY)); |
| 78 | ASSERT(!serverList.more && serverList.size() < CLIENT_KNOBS->TOO_MANY); |
| 79 | |
| 80 | for (auto s = serverList.begin(); s != serverList.end(); ++s) |
| 81 | servers.insert(decodeServerListValue(s->value).id()); |
| 82 | |
| 83 | ASSERT(servers.size() > 0); |
| 84 | } |
no test coverage detected