| 240 | }; |
| 241 | |
| 242 | class PaxosConfigTransactionImpl { |
| 243 | std::vector<ConfigTransactionInterface> ctis; |
| 244 | GetGenerationQuorum getGenerationQuorum; |
| 245 | CommitQuorum commitQuorum; |
| 246 | int numRetries{ 0 }; |
| 247 | Optional<UID> dID; |
| 248 | Database cx; |
| 249 | |
| 250 | ACTOR static Future<Optional<Value>> get(PaxosConfigTransactionImpl* self, Key key) { |
| 251 | state ConfigKey configKey = ConfigKey::decodeKey(key); |
| 252 | loop { |
| 253 | try { |
| 254 | state ConfigGeneration generation = wait(self->getGenerationQuorum.getGeneration()); |
| 255 | state std::vector<ConfigTransactionInterface> readReplicas = |
| 256 | self->getGenerationQuorum.getReadReplicas(); |
| 257 | std::vector<Future<Void>> fs; |
| 258 | for (ConfigTransactionInterface& readReplica : readReplicas) { |
| 259 | if (readReplica.hostname.present()) { |
| 260 | fs.push_back(tryInitializeRequestStream( |
| 261 | &readReplica.get, readReplica.hostname.get(), WLTOKEN_CONFIGTXN_GET)); |
| 262 | } |
| 263 | } |
| 264 | wait(waitForAll(fs)); |
| 265 | state Reference<ConfigTransactionInfo> configNodes(new ConfigTransactionInfo(readReplicas)); |
| 266 | ConfigTransactionGetReply reply = |
| 267 | wait(timeoutError(basicLoadBalance(configNodes, |
| 268 | &ConfigTransactionInterface::get, |
| 269 | ConfigTransactionGetRequest{ generation, configKey }), |
| 270 | CLIENT_KNOBS->GET_KNOB_TIMEOUT)); |
| 271 | if (reply.value.present()) { |
| 272 | return reply.value.get().toValue(); |
| 273 | } else { |
| 274 | return Optional<Value>{}; |
| 275 | } |
| 276 | } catch (Error& e) { |
| 277 | if (e.code() != error_code_timed_out && e.code() != error_code_broken_promise) { |
| 278 | throw; |
| 279 | } |
| 280 | self->reset(); |
| 281 | } |
| 282 | } |
| 283 | } |
| 284 | |
| 285 | ACTOR static Future<RangeResult> getConfigClasses(PaxosConfigTransactionImpl* self) { |
| 286 | state ConfigGeneration generation = wait(self->getGenerationQuorum.getGeneration()); |
| 287 | state std::vector<ConfigTransactionInterface> readReplicas = self->getGenerationQuorum.getReadReplicas(); |
| 288 | std::vector<Future<Void>> fs; |
| 289 | for (ConfigTransactionInterface& readReplica : readReplicas) { |
| 290 | if (readReplica.hostname.present()) { |
| 291 | fs.push_back(tryInitializeRequestStream( |
| 292 | &readReplica.getClasses, readReplica.hostname.get(), WLTOKEN_CONFIGTXN_GETCLASSES)); |
| 293 | } |
| 294 | } |
| 295 | wait(waitForAll(fs)); |
| 296 | state Reference<ConfigTransactionInfo> configNodes(new ConfigTransactionInfo(readReplicas)); |
| 297 | ConfigTransactionGetConfigClassesReply reply = |
| 298 | wait(basicLoadBalance(configNodes, |
| 299 | &ConfigTransactionInterface::getClasses, |