QueueRunner class imitates the behavior of the python version of QueueRunner which creates a thread for each enqueue op, runs close op on completion.
| 36 | /// QueueRunner class imitates the behavior of the python version of QueueRunner |
| 37 | /// which creates a thread for each enqueue op, runs close op on completion. |
| 38 | class QueueRunner : public RunnerInterface { |
| 39 | public: |
| 40 | /// Creates a new QueueRunner from proto. |
| 41 | // TODO(yuefengz): we may want to initialize from queues and ops in the |
| 42 | // future. |
| 43 | static Status New(const QueueRunnerDef& queue_runner_def, |
| 44 | std::unique_ptr<QueueRunner>* result); |
| 45 | |
| 46 | /// Creates a new QueueRunner with a coordinator, see coordinator.h for usage. |
| 47 | static Status New(const QueueRunnerDef& queue_runner_def, Coordinator* coord, |
| 48 | std::unique_ptr<QueueRunner>* result); |
| 49 | |
| 50 | /// Adds a callback that the queue runner will call when it detects an error. |
| 51 | void AddErrorCallback(const std::function<void(Status)>& cb); |
| 52 | |
| 53 | /// Delete the previously registered callbacks. |
| 54 | void ClearErrorCallbacks(); |
| 55 | |
| 56 | /// The destructor would join all the threads. |
| 57 | ~QueueRunner(); |
| 58 | |
| 59 | /// Starts the queue runner with the given session. |
| 60 | Status Start(Session* sess); |
| 61 | |
| 62 | /// Starts the queue runner with the given session and sets the run arguments |
| 63 | /// for sess->Run. It also collects and stores the cost model. |
| 64 | Status StartAndCollectCostGraph(Session* sess, |
| 65 | const RunOptions& run_options = RunOptions()); |
| 66 | |
| 67 | /// Starts the queue runner with the given session, and wait for up to the |
| 68 | /// specified time (in milliseconds) for the queues to start to fill up. |
| 69 | Status Start(Session* sess, int wait_for_ms); |
| 70 | Status StartAndCollectCostGraph(Session* session, int wait_for_ms, |
| 71 | const RunOptions& run_options = RunOptions()); |
| 72 | |
| 73 | /// Requests to stop and runs the cancel op. It would be called in a separate |
| 74 | /// thread when coordinator is set. If there is no coordinator it should be |
| 75 | /// called before calling Join. |
| 76 | void Stop(Session* sess); |
| 77 | |
| 78 | /// Joins all the threads. Returns okay if all threads run successfully; |
| 79 | /// otherwise returns the first captured failure status. |
| 80 | Status Join() final; |
| 81 | |
| 82 | /// Returns the latest status. |
| 83 | Status GetStatus(); |
| 84 | |
| 85 | // Returns the stored cost model. |
| 86 | Status ExportCostGraph(CostGraphDef* cost_graph) const override; |
| 87 | |
| 88 | private: |
| 89 | QueueRunner() : coord_(nullptr), stopped_(false), cg_mu_(nullptr) {} |
| 90 | |
| 91 | // Initializes the instance with the QueueRunnerDef proto. |
| 92 | Status Init(const QueueRunnerDef& queue_runner_def); |
| 93 | |
| 94 | // The Run function for each thread. |
| 95 | void Run(Session* sess, const string& enqueue_op); |