| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 536ebb9 commit a3b465a
7 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -47,8 +47,6 @@ struct IReactiveEngine | |||
| 47 | 47 | ||
| 48 | 48 | template <typename F> | |
| 49 | 49 | bool TryMerge(F&& f) { return false; } | |
| 50 | - | ||
| 51 | - void HintUpdateDuration(NodeT& node, uint dur) {} | ||
| 52 | 50 | }; | |
| 53 | 51 | ||
| 54 | 52 | /////////////////////////////////////////////////////////////////////////////////////////////////// | |
@@ -160,18 +158,13 @@ struct EngineInterface | |||
| 160 | 158 | { | |
| 161 | 159 | return Engine().TryMerge(std::forward<F>(f)); | |
| 162 | 160 | } | |
| 163 | - | ||
| 164 | - static void HintUpdateDuration(NodeT& node, uint dur) | ||
| 165 | - { | ||
| 166 | - Engine().HintUpdateDuration(node, dur); | ||
| 167 | - } | ||
| 168 | - | ||
| 169 | 161 | }; | |
| 170 | 162 | ||
| 171 | 163 | /////////////////////////////////////////////////////////////////////////////////////////////////// | |
| 172 | 164 | /// Engine traits | |
| 173 | 165 | /////////////////////////////////////////////////////////////////////////////////////////////////// | |
| 174 | 166 | template <typename> struct EnableNodeUpdateTimer : std::false_type {}; | |
| 175 | 167 | template <typename> struct EnableParallelUpdating : std::false_type {}; | |
| 168 | + template <typename> struct EnableConcurrentInput : std::false_type {}; | ||
| 176 | 169 | ||
| 177 | 170 | /****************************************/ REACT_IMPL_END /***************************************/ | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -85,7 +85,6 @@ class Node : public IReactiveNode | |||
| 85 | 85 | NodeVector<Node> Successors; | |
| 86 | 86 | ||
| 87 | 87 | ENodeState State = ENodeState::unchanged; | |
| 88 | - uint Weight = 0; | ||
| 89 | 88 | ||
| 90 | 89 | private: | |
| 91 | 90 | atomic<int> counter_ = 0; | |
@@ -114,8 +113,6 @@ class EngineBase : public IReactiveEngine<Node,TTurn> | |||
| 114 | 113 | void OnDynamicNodeAttach(Node& node, Node& parent, TTurn& turn); | |
| 115 | 114 | void OnDynamicNodeDetach(Node& node, Node& parent, TTurn& turn); | |
| 116 | 115 | ||
| 117 | - void HintUpdateDuration(Node& node, uint dur); | ||
| 118 | - | ||
| 119 | 116 | private: | |
| 120 | 117 | NodeVectT changedInputs_; | |
| 121 | 118 | empty_task* rootTask_ = new(task::allocate_root()) empty_task; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -52,10 +52,6 @@ class Node : public IReactiveNode | |||
| 52 | 52 | inline void SetQueuedFlag() { flags_.Set<flag_queued>(); } | |
| 53 | 53 | inline void ClearQueuedFlag() { flags_.Clear<flag_queued>(); } | |
| 54 | 54 | ||
| 55 | - inline bool IsHeavy() const { return flags_.Test<flag_heavy>(); } | ||
| 56 | - inline void SetHeavyFlag() { flags_.Set<flag_heavy>(); } | ||
| 57 | - inline void ClearHeavyFlag() { flags_.Clear<flag_heavy>(); } | ||
| 58 | - | ||
| 59 | 55 | inline bool IsMarked() const { return flags_.Test<flag_marked>(); } | |
| 60 | 56 | inline void SetMarkedFlag() { flags_.Set<flag_marked>(); } | |
| 61 | 57 | inline void ClearMarkedFlag() { flags_.Clear<flag_marked>(); } | |
@@ -105,7 +101,6 @@ class Node : public IReactiveNode | |||
| 105 | 101 | enum EFlags : uint16_t | |
| 106 | 102 | { | |
| 107 | 103 | flag_queued = 0, | |
| 108 | - flag_heavy, | ||
| 109 | 104 | flag_marked, | |
| 110 | 105 | flag_changed, | |
| 111 | 106 | flag_deferred, | |
@@ -142,8 +137,6 @@ class EngineBase : public IReactiveEngine<Node,TTurn> | |||
| 142 | 137 | void OnDynamicNodeAttach(Node& node, Node& parent, TTurn& turn); | |
| 143 | 138 | void OnDynamicNodeDetach(Node& node, Node& parent, TTurn& turn); | |
| 144 | 139 | ||
| 145 | - void HintUpdateDuration(Node& node, uint dur); | ||
| 146 | - | ||
| 147 | 140 | private: | |
| 148 | 141 | void applyAsyncDynamicAttach(Node& node, Node& parent, TTurn& turn); | |
| 149 | 142 | void applyAsyncDynamicDetach(Node& node, Node& parent, TTurn& turn); | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -76,7 +76,6 @@ class ParNode : public IReactiveNode | |||
| 76 | 76 | int Level = 0; | |
| 77 | 77 | int NewLevel = 0; | |
| 78 | 78 | atomic<bool> Collected = false; | |
| 79 | - uint Weight = 1; | ||
| 80 | 79 | ||
| 81 | 80 | NodeVector<ParNode> Successors; | |
| 82 | 81 | InvalidateMutexT InvalidateMutex; | |
@@ -156,8 +155,6 @@ class ParEngineBase : public EngineBase<ParNode,TTurn> | |||
| 156 | 155 | void OnDynamicNodeAttach(ParNode& node, ParNode& parent, TTurn& turn); | |
| 157 | 156 | void OnDynamicNodeDetach(ParNode& node, ParNode& parent, TTurn& turn); | |
| 158 | 157 | ||
| 159 | - void HintUpdateDuration(ParNode& node, uint dur); | ||
| 160 | - | ||
| 161 | 158 | private: | |
| 162 | 159 | void applyDynamicAttach(ParNode& node, ParNode& parent, TTurn& turn); | |
| 163 | 160 | void applyDynamicDetach(ParNode& node, ParNode& parent, TTurn& turn); | |
@@ -247,7 +244,7 @@ class PipeliningTurn : public TurnBase | |||
| 247 | 244 | PipeliningTurn* successor_ = nullptr; | |
| 248 | 245 | ||
| 249 | 246 | int currentLevel_ = -1; | |
| 250 | - int maxLevel_ = numeric_limits<int>::max(); /// This turn may only advance up to maxLevel | ||
| 247 | + int maxLevel_ = (numeric_limits<int>::max)(); /// This turn may only advance up to maxLevel | ||
| 251 | 248 | int minLevel_ = -1; /// successor.maxLevel = this.minLevel - 1 | |
| 252 | 249 | ||
| 253 | 250 | int curUpperBound_ = -1; | |
@@ -363,4 +360,9 @@ template <> struct EnableParallelUpdating<ToposortEngine<parallel>> : std::true_ | |||
| 363 | 360 | template <> struct EnableParallelUpdating<ToposortEngine<parallel_queue>> : std::true_type {}; | |
| 364 | 361 | template <> struct EnableParallelUpdating<ToposortEngine<parallel_pipeline>> : std::true_type {}; | |
| 365 | 362 | ||
| 363 | + template <typename> struct EnableConcurrentInput; | ||
| 364 | + template <> struct EnableConcurrentInput<ToposortEngine<sequential_queue>> : std::true_type {}; | ||
| 365 | + template <> struct EnableConcurrentInput<ToposortEngine<parallel_queue>> : std::true_type {}; | ||
| 366 | + template <> struct EnableConcurrentInput<ToposortEngine<parallel_pipeline>> : std::true_type {}; | ||
| 367 | + | ||
| 366 | 368 | /****************************************/ REACT_IMPL_END /***************************************/ | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -133,7 +133,7 @@ class UpdaterTask: public task | |||
| 133 | 133 | continue; | |
| 134 | 134 | ||
| 135 | 135 | // Heavyweight - spawn new task | |
| 136 | - if (succ->Weight > heavy_weight) | ||
| 136 | + if (succ->IsHeavyweight()) | ||
| 137 | 137 | { | |
| 138 | 138 | auto& t = *new(task::allocate_additional_child_of(*parent())) | |
| 139 | 139 | UpdaterTask(turn_, succ); | |
@@ -287,12 +287,6 @@ void EngineBase<TTurn>::OnDynamicNodeDetach(Node& node, Node& parent, TTurn& tur | |||
| 287 | 287 | parent.Successors.Remove(node); | |
| 288 | 288 | }// ~parent.ShiftMutex (write) | |
| 289 | 289 | ||
| 290 | - template <typename TTurn> | ||
| 291 | - void EngineBase<TTurn>::HintUpdateDuration(Node& node, uint dur) | ||
| 292 | - { | ||
| 293 | - node.Weight = dur; | ||
| 294 | - } | ||
| 295 | - | ||
| 296 | 290 | // Explicit instantiation | |
| 297 | 291 | template class EngineBase<Turn>; | |
| 298 | 292 | template class EngineBase<DefaultQueueableTurn<Turn>>; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -112,7 +112,7 @@ class UpdaterTask: public task | |||
| 112 | 112 | succ->SetReadyCount(0); | |
| 113 | 113 | ||
| 114 | 114 | // Heavyweight - spawn new task | |
| 115 | - if (succ->IsHeavy()) | ||
| 115 | + if (succ->IsHeavyweight()) | ||
| 116 | 116 | { | |
| 117 | 117 | auto& t = *new(task::allocate_additional_child_of(*parent())) | |
| 118 | 118 | UpdaterTask(turn_, succ); | |
@@ -155,7 +155,7 @@ void EngineBase<TTurn>::OnTurnPropagate(TTurn& turn) | |||
| 155 | 155 | // Phase 1 | |
| 156 | 156 | while (scheduledNodes_.FetchNext()) | |
| 157 | 157 | { | |
| 158 | - for (auto* curNode : scheduledNodes_.NextNodes()) | ||
| 158 | + for (auto* curNode : scheduledNodes_.NextValues()) | ||
| 159 | 159 | { | |
| 160 | 160 | if (curNode->Level < curNode->NewLevel) | |
| 161 | 161 | { | |
@@ -244,15 +244,6 @@ void EngineBase<TTurn>::OnDynamicNodeDetach(Node& node, Node& parent, TTurn& tur | |||
| 244 | 244 | OnNodeDetach(node, parent); | |
| 245 | 245 | } | |
| 246 | 246 | ||
| 247 | - template <typename TTurn> | ||
| 248 | - void EngineBase<TTurn>::HintUpdateDuration(Node& node, uint dur) | ||
| 249 | - { | ||
| 250 | - if (dur > heavy_weight) | ||
| 251 | - node.SetHeavyFlag(); | ||
| 252 | - else | ||
| 253 | - node.ClearHeavyFlag(); | ||
| 254 | - } | ||
| 255 | - | ||
| 256 | 247 | template <typename TTurn> | |
| 257 | 248 | void EngineBase<TTurn>::applyAsyncDynamicAttach(Node& node, Node& parent, TTurn& turn) | |
| 258 | 249 | { | |
@@ -306,7 +297,7 @@ void EngineBase<TTurn>::processChildren(Node& node, TTurn& turn) | |||
| 306 | 297 | continue; | |
| 307 | 298 | ||
| 308 | 299 | // Light nodes use sequential toposort in phase 1 | |
| 309 | - if (! succ->IsHeavy()) | ||
| 300 | + if (! succ->IsHeavyweight()) | ||
| 310 | 301 | { | |
| 311 | 302 | if (!succ->IsQueued()) | |
| 312 | 303 | { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -62,7 +62,7 @@ void SeqEngineBase<TTurn>::OnTurnPropagate(TTurn& turn) | |||
| 62 | 62 | { | |
| 63 | 63 | while (scheduledNodes_.FetchNext()) | |
| 64 | 64 | { | |
| 65 | - for (auto* curNode : scheduledNodes_.NextNodes()) | ||
| 65 | + for (auto* curNode : scheduledNodes_.NextValues()) | ||
| 66 | 66 | { | |
| 67 | 67 | if (curNode->Level < curNode->NewLevel) | |
| 68 | 68 | { | |
@@ -133,15 +133,17 @@ void ParEngineBase<TTurn>::OnTurnPropagate(TTurn& turn) | |||
| 133 | 133 | while (topoQueue_.FetchNext()) | |
| 134 | 134 | { | |
| 135 | 135 | //using RangeT = tbb::blocked_range<vector<ParNode*>::const_iterator>; | |
| 136 | - using RangeT = ParEngineBase::TopoQueueT::RangeT; | ||
| 136 | + using RangeT = ParEngineBase::TopoQueueT::NextRangeT; | ||
| 137 | 137 | ||
| 138 | 138 | // Iterate all nodes of current level and start processing them in parallel | |
| 139 | 139 | tbb::parallel_for( | |
| 140 | 140 | topoQueue_.NextRange(), | |
| 141 | 141 | [&] (const RangeT& range) | |
| 142 | 142 | { | |
| 143 | - for (auto* curNode : range) | ||
| 143 | + for (const auto& e : range) | ||
| 144 | 144 | { | |
| 145 | + auto* curNode = e.first; | ||
| 146 | + | ||
| 145 | 147 | if (curNode->Level < curNode->NewLevel) | |
| 146 | 148 | { | |
| 147 | 149 | curNode->Level = curNode->NewLevel; | |
@@ -186,15 +188,6 @@ void ParEngineBase<TTurn>::OnDynamicNodeDetach(ParNode& node, ParNode& parent, T | |||
| 186 | 188 | dynRequests_.push_back(data); | |
| 187 | 189 | } | |
| 188 | 190 | ||
| 189 | - template <typename TTurn> | ||
| 190 | - void ParEngineBase<TTurn>::HintUpdateDuration(ParNode& node, uint dur) | ||
| 191 | - { | ||
| 192 | - if (dur < min_weight) | ||
| 193 | - dur = min_weight; | ||
| 194 | - | ||
| 195 | - node.Weight = dur; | ||
| 196 | - } | ||
| 197 | - | ||
| 198 | 191 | template <typename TTurn> | |
| 199 | 192 | void ParEngineBase<TTurn>::applyDynamicAttach(ParNode& node, ParNode& parent, TTurn& turn) | |
| 200 | 193 | { | |
@@ -331,7 +324,7 @@ void PipeliningTurn::Remove() | |||
| 331 | 324 | } | |
| 332 | 325 | else if (successor_) | |
| 333 | 326 | { | |
| 334 | - successor_->SetMaxLevel(numeric_limits<int>::max()); | ||
| 327 | + successor_->SetMaxLevel((numeric_limits<int>::max)()); | ||
| 335 | 328 | successor_->predecessor_ = nullptr; | |
| 336 | 329 | } | |
| 337 | 330 | ||
@@ -434,10 +427,10 @@ void PipeliningEngine::OnTurnPropagate(PipeliningTurn& turn) | |||
| 434 | 427 | ||
| 435 | 428 | while (turn.TopoQueue.FetchNext()) | |
| 436 | 429 | { | |
| 437 | - using RangeT = PipeliningTurn::TopoQueueT::RangeT; | ||
| 430 | + using RangeT = PipeliningTurn::TopoQueueT::NextRangeT; | ||
| 438 | 431 | ||
| 439 | - for (const auto* node : turn.TopoQueue.NextRange()) | ||
| 440 | - turn.AdjustUpperBound(node->Level); | ||
| 432 | + for (const auto& e : turn.TopoQueue.NextRange()) | ||
| 433 | + turn.AdjustUpperBound(e.first->Level); | ||
| 441 | 434 | ||
| 442 | 435 | advanceTurn(turn); | |
| 443 | 436 | ||
@@ -446,8 +439,9 @@ void PipeliningEngine::OnTurnPropagate(PipeliningTurn& turn) | |||
| 446 | 439 | turn.TopoQueue.NextRange(), | |
| 447 | 440 | [&] (const RangeT& range) | |
| 448 | 441 | { | |
| 449 | - for (auto* curNode : range) | ||
| 442 | + for (const auto& e : range) | ||
| 450 | 443 | { | |
| 444 | + auto* curNode = e.first; | ||
| 451 | 445 | if (curNode->Level < curNode->NewLevel) | |
| 452 | 446 | { | |
| 453 | 447 | curNode->Level = curNode->NewLevel; | |
@@ -492,7 +486,7 @@ void PipeliningEngine::OnDynamicNodeDetach(ParNode& node, ParNode& parent, Pipel | |||
| 492 | 486 | ||
| 493 | 487 | void PipeliningEngine::applyDynamicAttach(ParNode& node, ParNode& parent, PipeliningTurn& turn) | |
| 494 | 488 | { | |
| 495 | - turn.WaitForMaxLevel(numeric_limits<int>::max()); | ||
| 489 | + turn.WaitForMaxLevel((numeric_limits<int>::max)()); | ||
| 496 | 490 | ||
| 497 | 491 | OnNodeAttach(node, parent); | |
| 498 | 492 | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments