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

Function ThreadRun

serving/processor/storage/feature_store_mgr.cc:25–95  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

23
24namespace {
25void ThreadRun(AsyncFeatureStoreMgr* mgr, int idx,
26 bool is_update_thread) {
27 std::mutex* mu = nullptr;
28 std::condition_variable* cv = nullptr;
29 sparse_task_queue* queue = nullptr;
30
31 if (is_update_thread) {
32 mu = mgr->GetUpdateMutex(idx);
33 cv = mgr->GetUpdateCV(idx);
34 queue = mgr->GetUpdateSparseTaskQueue(idx);
35 } else {
36 mu = mgr->GetMutex(idx);
37 cv = mgr->GetCV(idx);
38 queue = mgr->GetSparseTaskQueue(idx);
39 }
40
41 const int try_count = 64;
42 int curr_try_count = 0;
43 SparseTask* task = nullptr;
44 bool succeeded = false;
45
46 while ((succeeded = queue->try_dequeue(task)) ||
47 !mgr->ShouldStop()) {
48 if (!succeeded) {
49 ++curr_try_count;
50 if (curr_try_count <= try_count) {
51 continue;
52 }
53 curr_try_count = 0;
54
55 if (is_update_thread) {
56 // Ready going to sleep
57 *(mgr->GetUpdateSleepingFlag(idx)) = true;
58 *(mgr->GetUpdateReadyFlag(idx)) = false;
59 } else {
60 // Ready going to sleep
61 *(mgr->GetSleepingFlag(idx)) = true;
62 *(mgr->GetReadyFlag(idx)) = false;
63 }
64
65 {
66 // try to wait signal when have no elements in the queue
67 std::unique_lock<std::mutex> lock(*mu);
68 cv->wait(lock, [is_update_thread, mgr, idx] {
69 return (is_update_thread ?
70 *(mgr->GetUpdateReadyFlag(idx)) :
71 *(mgr->GetReadyFlag(idx))) ||
72 mgr->ShouldStop();
73 });
74 lock.unlock();
75 }
76
77 if (is_update_thread) {
78 *(mgr->GetUpdateSleepingFlag(idx)) = false;
79 } else {
80 *(mgr->GetSleepingFlag(idx)) = false;
81 }
82

Callers

nothing calls this directly

Calls 14

GetUpdateMutexMethod · 0.80
GetUpdateCVMethod · 0.80
GetMutexMethod · 0.80
GetCVMethod · 0.80
GetSparseTaskQueueMethod · 0.80
GetUpdateSleepingFlagMethod · 0.80
GetUpdateReadyFlagMethod · 0.80
GetSleepingFlagMethod · 0.80
GetReadyFlagMethod · 0.80
ShouldStopMethod · 0.45
waitMethod · 0.45

Tested by

no test coverage detected