| 45 | }; |
| 46 | |
| 47 | void RdmaRemoteRendezvous::RecvFromRemoteAsync( |
| 48 | const Rendezvous::ParsedKey& parsed, const Rendezvous::Args& recv_args, |
| 49 | DoneCallback done) { |
| 50 | Status s; |
| 51 | // parse src_name and dst_name |
| 52 | string src_name, dst_name, unused; |
| 53 | if (!DeviceNameUtils::SplitDeviceName(parsed.src_device, &src_name, |
| 54 | &unused) || |
| 55 | !DeviceNameUtils::SplitDeviceName(parsed.dst_device, &dst_name, |
| 56 | &unused)) { |
| 57 | s = errors::Internal("Could not parse src or dst name."); |
| 58 | } |
| 59 | if (!s.ok()) { |
| 60 | LOG(ERROR) << "s is not ok, error code " << s.error_message(); |
| 61 | done(s, Args(), recv_args, Tensor{}, false); |
| 62 | return; |
| 63 | } |
| 64 | CHECK(dst_name.compare(rdma_mgr_->local_worker()) == 0); |
| 65 | RdmaChannel* rc = rdma_mgr_->FindChannel(src_name); |
| 66 | string key(parsed.FullKey()); |
| 67 | string key_with_step_id = VerbsUtil::AppendStepidToKey(key, step_id_); |
| 68 | |
| 69 | Device* dst_dev; |
| 70 | s = env_->device_mgr->LookupDevice(parsed.dst_device, &dst_dev); |
| 71 | CHECK(s.ok()) << "s is not ok, error code " << s.error_message(); |
| 72 | if (!s.ok()) { |
| 73 | done(s, Args(), recv_args, Tensor(), true); |
| 74 | return; |
| 75 | } |
| 76 | |
| 77 | RdmaTensorRequest* request = |
| 78 | rc->InsertTensorRequest(key, step_id_, dst_dev, recv_args, done); |
| 79 | request->Start(); |
| 80 | } |
| 81 | |
| 82 | RdmaRendezvousMgr::RdmaRendezvousMgr(const WorkerEnv* env) |
| 83 | : BaseRendezvousMgr(env) {} |
nothing calls this directly
no test coverage detected