MCPcopy Create free account
hub / github.com/DeepRec-AI/DeepRec / RecvFromRemoteAsync

Method RecvFromRemoteAsync

tensorflow/contrib/verbs/rdma_rendezvous_mgr.cc:47–80  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

45};
46
47void 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
82RdmaRendezvousMgr::RdmaRendezvousMgr(const WorkerEnv* env)
83 : BaseRendezvousMgr(env) {}

Callers

nothing calls this directly

Calls 10

InternalFunction · 0.85
FindChannelMethod · 0.80
FullKeyMethod · 0.80
LookupDeviceMethod · 0.80
InsertTensorRequestMethod · 0.80
ArgsClass · 0.50
TensorClass · 0.50
okMethod · 0.45
compareMethod · 0.45
StartMethod · 0.45

Tested by

no test coverage detected