| 219 | } |
| 220 | |
| 221 | Status BaseRemoteRendezvous::Initialize(WorkerSession* session) { |
| 222 | CHECK_NE(session, nullptr) << "session must not be null!"; |
| 223 | std::vector<DeferredCall> deferred_calls; |
| 224 | { |
| 225 | mutex_lock l(mu_); |
| 226 | if (session_ != nullptr) { |
| 227 | if (worker_name_ != session->worker_name) { |
| 228 | Status s = errors::Internal( |
| 229 | "Double init! Worker names would have changed from: ", |
| 230 | worker_name_, " -> ", session->worker_name); |
| 231 | LOG(WARNING) << s; |
| 232 | return s; |
| 233 | } |
| 234 | } |
| 235 | session_ = session; |
| 236 | worker_name_ = session->worker_name; |
| 237 | std::swap(deferred_calls, deferred_calls_); |
| 238 | } |
| 239 | for (auto& call : deferred_calls) { |
| 240 | RecvLocalAsyncInternal(call.parsed, std::move(call.done)); |
| 241 | } |
| 242 | |
| 243 | std::vector<DeferredFuseCall> deferred_fuse_calls; |
| 244 | { |
| 245 | mutex_lock l(mu_); |
| 246 | std::swap(deferred_fuse_calls, deferred_fuse_calls_); |
| 247 | } |
| 248 | for (auto& fuse_call : deferred_fuse_calls) { |
| 249 | FuseRecvLocalAsyncInternal(fuse_call.parsed_keys, |
| 250 | std::move(fuse_call.done)); |
| 251 | } |
| 252 | |
| 253 | std::vector<DeferredFlowControlCall> deferred_flow_control_calls; |
| 254 | { |
| 255 | mutex_lock l(mu_); |
| 256 | std::swap(deferred_flow_control_calls, deferred_flow_control_calls_); |
| 257 | } |
| 258 | for (auto& fc_call : deferred_flow_control_calls) { |
| 259 | FlowControlRecvLocalAsyncInternal(fc_call.tag, fc_call.parsed, |
| 260 | std::move(fc_call.done)); |
| 261 | } |
| 262 | |
| 263 | return Status::OK(); |
| 264 | } |
| 265 | |
| 266 | WorkerSession* BaseRemoteRendezvous::session() { |
| 267 | tf_shared_lock l(mu_); |