MCPcopy Create free account
hub / github.com/ml-explore/mlx / worker

Method worker

mlx/distributed/ring/ring.cpp:202–266  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

200 }
201
202 void worker() {
203 int error_count = 0;
204 bool delete_recv = false;
205 bool delete_send = false;
206 while (true) {
207 {
208 std::unique_lock lock(queue_mutex_);
209
210 if (delete_recv) {
211 recvs_.front().promise.set_value();
212 recvs_.pop_front();
213 delete_recv = false;
214 }
215 if (delete_send) {
216 sends_.front().promise.set_value();
217 sends_.pop_front();
218 delete_send = false;
219 }
220
221 if (stop_) {
222 return;
223 }
224
225 if (!have_tasks()) {
226 condition_.wait(lock, [this] { return stop_ || have_tasks(); });
227 if (stop_) {
228 return;
229 }
230 }
231 }
232
233 if (!recvs_.empty()) {
234 auto& task = recvs_.front();
235 ssize_t r = ::recv(fd_, task.buffer, task.size, 0);
236 if (r > 0) {
237 task.buffer = static_cast<char*>(task.buffer) + r;
238 task.size -= r;
239 delete_recv = task.size == 0;
240 error_count = 0;
241 } else if (errno != EAGAIN) {
242 error_count++;
243 log_info(
244 true, "Receiving from socket", fd_, "failed with errno", errno);
245 }
246 }
247 if (!sends_.empty()) {
248 auto& task = sends_.front();
249 ssize_t r = ::send(fd_, task.buffer, task.size, 0);
250 if (r > 0) {
251 task.buffer = static_cast<char*>(task.buffer) + r;
252 task.size -= r;
253 delete_send = task.size == 0;
254 error_count = 0;
255 } else if (errno != EAGAIN) {
256 error_count++;
257 log_info(true, "Sending to socket", fd_, "failed with errno", errno);
258 }
259 }

Callers

nothing calls this directly

Calls 6

log_infoFunction · 0.85
set_valueMethod · 0.80
recvFunction · 0.50
sendFunction · 0.50
waitMethod · 0.45
emptyMethod · 0.45

Tested by

no test coverage detected