Cpp-Taskflow  2.2.0
observer.hpp
1 // 2019/07/31 - modified by Tsung-Wei Huang
2 // - fixed the missing comma in outputing JSON
3 //
4 // 2019/06/13 - modified by Tsung-Wei Huang
5 // - added TaskView interface
6 //
7 // 2019/04/17 - created by Tsung-Wei Huang
8 
9 #pragma once
10 
11 #include <iostream>
12 #include <sstream>
13 #include <vector>
14 #include <cstdlib>
15 #include <cstdio>
16 #include <atomic>
17 #include <memory>
18 #include <deque>
19 #include <optional>
20 #include <thread>
21 #include <algorithm>
22 #include <set>
23 #include <numeric>
24 #include <cassert>
25 
26 #include "task.hpp"
27 
28 namespace tf {
29 
40 
41  public:
42 
46  virtual ~ExecutorObserverInterface() = default;
47 
52  virtual void set_up(unsigned num_workers) = 0;
53 
59  virtual void on_entry(unsigned worker_id, TaskView task_view) = 0;
60 
66  virtual void on_exit(unsigned worker_id, TaskView task_view) = 0;
67 };
68 
69 // ------------------------------------------------------------------
70 
78 
79  friend class Executor;
80 
81  // data structure to record each task execution
82  struct Execution {
83 
84  TaskView task_view;
85 
88 
89  Execution(
90  TaskView tv,
92  ) :
93  task_view {tv}, beg {b} {
94  }
95 
96  Execution(
97  TaskView tv,
100  ) :
101  task_view {tv}, beg {b}, end {e} {
102  }
103  };
104 
105  // data structure to store the entire execution timeline
106  struct Timeline {
109  };
110 
111  public:
112 
117  inline void dump(std::ostream& ostream) const;
118 
123  inline std::string dump() const;
124 
128  inline void clear();
129 
134  inline size_t num_tasks() const;
135 
136  private:
137 
138  inline void set_up(unsigned num_workers) override final;
139  inline void on_entry(unsigned worker_id, TaskView task_view) override final;
140  inline void on_exit(unsigned worker_id, TaskView task_view) override final;
141 
142  Timeline _timeline;
143 };
144 
145 // Procedure: set_up
146 inline void ExecutorObserver::set_up(unsigned num_workers) {
147 
148  _timeline.executions.resize(num_workers);
149 
150  for(unsigned w=0; w<num_workers; ++w) {
151  _timeline.executions[w].reserve(1024);
152  }
153 
154  _timeline.origin = std::chrono::steady_clock::now();
155 }
156 
157 // Procedure: on_entry
158 inline void ExecutorObserver::on_entry(unsigned w, TaskView tv) {
159  _timeline.executions[w].emplace_back(tv, std::chrono::steady_clock::now());
160 }
161 
162 // Procedure: on_exit
163 inline void ExecutorObserver::on_exit(unsigned w, TaskView tv) {
164  static_cast<void>(tv); // avoid warning from compiler
165  assert(_timeline.executions[w].size() > 0);
166  _timeline.executions[w].back().end = std::chrono::steady_clock::now();
167 }
168 
169 // Function: clear
170 inline void ExecutorObserver::clear() {
171  for(size_t w=0; w<_timeline.executions.size(); ++w) {
172  _timeline.executions[w].clear();
173  }
174 }
175 
176 // Procedure: dump
177 inline void ExecutorObserver::dump(std::ostream& os) const {
178 
179  size_t first;
180 
181  for(first = 0; first<_timeline.executions.size(); ++first) {
182  if(_timeline.executions[first].size() > 0) {
183  break;
184  }
185  }
186 
187  os << '[';
188 
189  for(size_t w=first; w<_timeline.executions.size(); w++) {
190 
191  if(w != first && _timeline.executions[w].size() > 0) {
192  os << ',';
193  }
194 
195  for(size_t i=0; i<_timeline.executions[w].size(); i++) {
196 
197  os << '{'
198  << "\"cat\":\"ExecutorObserver\","
199  << "\"name\":\"" << _timeline.executions[w][i].task_view.name() << "\","
200  << "\"ph\":\"X\","
201  << "\"pid\":1,"
202  << "\"tid\":" << w << ','
204  _timeline.executions[w][i].beg - _timeline.origin
205  ).count() << ','
207  _timeline.executions[w][i].end - _timeline.executions[w][i].beg
208  ).count();
209 
210  if(i != _timeline.executions[w].size() - 1) {
211  os << "},";
212  }
213  else {
214  os << '}';
215  }
216  }
217  }
218  os << "]\n";
219 }
220 
221 // Function: dump
223  std::ostringstream oss;
224  dump(oss);
225  return oss.str();
226 }
227 
228 // Function: num_tasks
229 inline size_t ExecutorObserver::num_tasks() const {
230  return std::accumulate(
231  _timeline.executions.begin(), _timeline.executions.end(), size_t{0},
232  [](size_t sum, const auto& exe){
233  return sum + exe.size();
234  }
235  );
236 }
237 
238 
239 } // end of namespace tf -------------------------------------------
240 
241 
virtual void on_entry(unsigned worker_id, TaskView task_view)=0
method to call before a worker thread executes a closure
virtual void set_up(unsigned num_workers)=0
constructor-like method to call when the executor observer is fully created
Default executor observer to dump the execution timelines.
Definition: observer.hpp:77
virtual ~ExecutorObserverInterface()=default
virtual destructor
size_t num_tasks() const
get the number of total tasks in the observer
Definition: observer.hpp:229
T duration_cast(T... args)
Definition: taskflow.hpp:5
A constant wrapper class to a task node, mainly used in the tf::ExecutorObserver interface.
Definition: task.hpp:370
virtual void on_exit(unsigned worker_id, TaskView task_view)=0
method to call after a worker thread executed a closure
void clear()
clear the timeline data
Definition: observer.hpp:170
std::string dump() const
dump the timelines in JSON to a std::string
Definition: observer.hpp:222
The executor class to run a taskflow graph.
Definition: executor.hpp:88
The interface class for creating an executor observer.
Definition: observer.hpp:39