| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 234de7b commit 6bc05be
3 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -13,7 +13,13 @@ namespace tf { | |||
| 13 | 13 | ||
| 14 | 14 | // Procedure: _schedule_async_task | |
| 15 | 15 | TF_FORCE_INLINE void Executor::_schedule_async_task(Node* node) { | |
| 16 | - (pt::this_worker) ? _schedule(*pt::this_worker, node) : _schedule(node); | ||
| 16 | + auto w = this_worker(); | ||
| 17 | + if(w) { | ||
| 18 | + _schedule(*w, node); | ||
| 19 | + } | ||
| 20 | + else{ | ||
| 21 | + _schedule(node); | ||
| 22 | + } | ||
| 17 | 23 | } | |
| 18 | 24 | ||
| 19 | 25 | // Procedure: _tear_down_async | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -542,6 +542,16 @@ class Executor { | |||
| 542 | 542 | @endcode | |
| 543 | 543 | */ | |
| 544 | 544 | size_t num_taskflows() const; | |
| 545 | + | ||
| 546 | + /** | ||
| 547 | + @brief queries pointer to the current worker if it belongs to this executor, otherwise returns nullptr | ||
| 548 | + | ||
| 549 | + @code{.cpp} | ||
| 550 | + auto w = executor.this_worker(); | ||
| 551 | + assert(w == nullptr || w->executor() == &executor); | ||
| 552 | + @endcode | ||
| 553 | + */ | ||
| 554 | + Worker* this_worker(); | ||
| 545 | 555 | ||
| 546 | 556 | /** | |
| 547 | 557 | @brief queries the id of the caller thread within this executor | |
@@ -1063,6 +1073,8 @@ class Executor { | |||
| 1063 | 1073 | ||
| 1064 | 1074 | std::shared_ptr<WorkerInterface> _worker_interface; | |
| 1065 | 1075 | std::unordered_set<std::shared_ptr<ObserverInterface>> _observers; | |
| 1076 | + std::unordered_map<std::thread::id, size_t> _wids; | ||
| 1077 | + | ||
| 1066 | 1078 | ||
| 1067 | 1079 | void _shutdown(); | |
| 1068 | 1080 | void _observer_prologue(Worker&, Node*); | |
@@ -1224,23 +1236,31 @@ inline size_t Executor::num_taskflows() const { | |||
| 1224 | 1236 | return _taskflows.size(); | |
| 1225 | 1237 | } | |
| 1226 | 1238 | ||
| 1239 | + inline Worker* Executor::this_worker() { | ||
| 1240 | + auto itr = _wids.find(std::this_thread::get_id()); | ||
| 1241 | + return itr == _wids.end() ? nullptr : &_workers[itr->second]; | ||
| 1242 | + } | ||
| 1243 | + | ||
| 1227 | 1244 | // Function: this_worker_id | |
| 1228 | 1245 | inline int Executor::this_worker_id() const { | |
| 1229 | - auto w = pt::this_worker; | ||
| 1230 | - return (w && w->_executor == this) ? static_cast<int>(w->_id) : -1; | ||
| 1246 | + auto i = _wids.find(std::this_thread::get_id()); | ||
| 1247 | + return i == _wids.end() ? -1 : static_cast<int>(_workers[i->second]._id); | ||
| 1231 | 1248 | } | |
| 1232 | 1249 | ||
| 1233 | 1250 | // Procedure: _spawn | |
| 1234 | 1251 | inline void Executor::_spawn(size_t N) { | |
| 1235 | - | ||
| 1252 | + std::mutex mutex; | ||
| 1236 | 1253 | for(size_t id=0; id<N; ++id) { | |
| 1237 | 1254 | _workers[id]._id = id; | |
| 1238 | 1255 | _workers[id]._vtm = id; | |
| 1239 | 1256 | _workers[id]._executor = this; | |
| 1240 | 1257 | _workers[id]._waiter = &_notifier._waiters[id]; | |
| 1241 | 1258 | _workers[id]._thread = std::thread([&, &w=_workers[id]] () { | |
| 1242 | 1259 | ||
| 1243 | - pt::this_worker = &w; | ||
| 1260 | + { | ||
| 1261 | + std::scoped_lock lock(mutex); | ||
| 1262 | + _wids[std::this_thread::get_id()] = w._id; | ||
| 1263 | + } | ||
| 1244 | 1264 | ||
| 1245 | 1265 | // initialize the random engine and seed for work-stealing loop | |
| 1246 | 1266 | w._rdgen.seed(static_cast<std::default_random_engine::result_type>( | |
@@ -2080,7 +2100,7 @@ tf::Future<void> Executor::run_until(Taskflow& f, P&& p, C&& c) { | |||
| 2080 | 2100 | std::lock_guard<std::mutex> lock(f._mutex); | |
| 2081 | 2101 | f._topologies.push(t); | |
| 2082 | 2102 | if(f._topologies.size() == 1) { | |
| 2083 | - _set_up_topology(pt::this_worker, t.get()); | ||
| 2103 | + _set_up_topology(this_worker(), t.get()); | ||
| 2084 | 2104 | } | |
| 2085 | 2105 | } | |
| 2086 | 2106 | ||
@@ -2108,23 +2128,25 @@ void Executor::corun(T& target) { | |||
| 2108 | 2128 | ||
| 2109 | 2129 | static_assert(has_graph_v<T>, "target must define a member function 'Graph& graph()'"); | |
| 2110 | 2130 | ||
| 2111 | - if(pt::this_worker == nullptr || pt::this_worker->_executor != this) { | ||
| 2131 | + Worker* w = this_worker(); | ||
| 2132 | + if(w == nullptr || w->_executor != this) { | ||
| 2112 | 2133 | TF_THROW("corun must be called by a worker of the executor"); | |
| 2113 | 2134 | } | |
| 2114 | 2135 | ||
| 2115 | 2136 | Node anchor; | |
| 2116 | - _corun_graph(*pt::this_worker, &anchor, target.graph().begin(), target.graph().end()); | ||
| 2137 | + _corun_graph(*w, &anchor, target.graph().begin(), target.graph().end()); | ||
| 2117 | 2138 | } | |
| 2118 | 2139 | ||
| 2119 | 2140 | // Function: corun_until | |
| 2120 | 2141 | template <typename P> | |
| 2121 | 2142 | void Executor::corun_until(P&& predicate) { | |
| 2122 | 2143 | ||
| 2123 | - if(pt::this_worker == nullptr || pt::this_worker->_executor != this) { | ||
| 2144 | + Worker* w = this_worker(); | ||
| 2145 | + if(w == nullptr || w->_executor != this) { | ||
| 2124 | 2146 | TF_THROW("corun_until must be called by a worker of the executor"); | |
| 2125 | 2147 | } | |
| 2126 | 2148 | ||
| 2127 | - _corun_until(*pt::this_worker, std::forward<P>(predicate)); | ||
| 2149 | + _corun_until(*w, std::forward<P>(predicate)); | ||
| 2128 | 2150 | } | |
| 2129 | 2151 | ||
| 2130 | 2152 | // Procedure: _corun_graph | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -106,20 +106,6 @@ class Worker { | |||
| 106 | 106 | BoundedTaskQueue<Node*> _wsq; | |
| 107 | 107 | }; | |
| 108 | 108 | ||
| 109 | - | ||
| 110 | - // ---------------------------------------------------------------------------- | ||
| 111 | - // Per-thread | ||
| 112 | - // ---------------------------------------------------------------------------- | ||
| 113 | - | ||
| 114 | - namespace pt { | ||
| 115 | - | ||
| 116 | - /** | ||
| 117 | - @private | ||
| 118 | - */ | ||
| 119 | - inline thread_local Worker* this_worker {nullptr}; | ||
| 120 | - | ||
| 121 | - } | ||
| 122 | - | ||
| 123 | 109 | // ---------------------------------------------------------------------------- | |
| 124 | 110 | // Class Definition: WorkerView | |
| 125 | 111 | // ---------------------------------------------------------------------------- | |
| Back | FazBrowse Home | New Git URL |
0 commit comments