| 465 | } |
| 466 | |
| 467 | std::shared_ptr<Transaction> RawSiteToSiteClient::createTransaction(TransferDirection direction) { |
| 468 | int ret; |
| 469 | bool dataAvailable; |
| 470 | std::shared_ptr<Transaction> transaction = nullptr; |
| 471 | |
| 472 | if (peer_state_ != READY) { |
| 473 | bootstrap(); |
| 474 | } |
| 475 | |
| 476 | if (peer_state_ != READY) { |
| 477 | return transaction; |
| 478 | } |
| 479 | |
| 480 | if (direction == RECEIVE) { |
| 481 | ret = writeRequestType(RECEIVE_FLOWFILES); |
| 482 | |
| 483 | if (ret <= 0) { |
| 484 | return transaction; |
| 485 | } |
| 486 | |
| 487 | RespondCode code; |
| 488 | std::string message; |
| 489 | |
| 490 | ret = readRespond(nullptr, code, message); |
| 491 | |
| 492 | if (ret <= 0) { |
| 493 | return transaction; |
| 494 | } |
| 495 | |
| 496 | org::apache::nifi::minifi::io::CRCStream<SiteToSitePeer> crcstream(gsl::make_not_null(peer_.get())); |
| 497 | switch (code) { |
| 498 | case MORE_DATA: |
| 499 | dataAvailable = true; |
| 500 | logger_->log_trace("Site2Site peer indicates that data is available"); |
| 501 | transaction = std::make_shared<Transaction>(direction, std::move(crcstream)); |
| 502 | known_transactions_[transaction->getUUID()] = transaction; |
| 503 | transaction->setDataAvailable(dataAvailable); |
| 504 | logger_->log_trace("Site2Site create transaction %s", transaction->getUUIDStr()); |
| 505 | return transaction; |
| 506 | case NO_MORE_DATA: |
| 507 | dataAvailable = false; |
| 508 | logger_->log_trace("Site2Site peer indicates that no data is available"); |
| 509 | transaction = std::make_shared<Transaction>(direction, std::move(crcstream)); |
| 510 | known_transactions_[transaction->getUUID()] = transaction; |
| 511 | transaction->setDataAvailable(dataAvailable); |
| 512 | logger_->log_trace("Site2Site create transaction %s", transaction->getUUIDStr()); |
| 513 | return transaction; |
| 514 | default: |
| 515 | logger_->log_warn("Site2Site got unexpected response %d when asking for data", code); |
| 516 | return NULL; |
| 517 | } |
| 518 | } else { |
| 519 | ret = writeRequestType(SEND_FLOWFILES); |
| 520 | |
| 521 | if (ret <= 0) { |
| 522 | return NULL; |
| 523 | } else { |
| 524 | org::apache::nifi::minifi::io::CRCStream<SiteToSitePeer> crcstream(gsl::make_not_null(peer_.get())); |
no test coverage detected