Line data Source code
1 : /*
2 : * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
3 : */
4 :
5 : #ifndef __WORK_PIPELINE_INL_H__
6 : #define __WORK_PIPELINE_INL_H__
7 :
8 :
9 : template<typename InputT, typename ResultT, typename ExternalT, typename SubResultT>
10 2064 : WorkStage<InputT, ResultT, ExternalT, SubResultT>::WorkStage(
11 : std::vector<std::pair<int,int> > tinfo,
12 : ExecuteFn efn, MergeFn mfn, int tid, int tinst) :
13 4128 : stage_(std::numeric_limits<uint32_t>::max()),
14 2064 : merger_(mfn),
15 2064 : efn_(efn),
16 2064 : finished_(false),
17 2064 : running_(false),
18 2064 : tid_(tid),
19 2064 : tinst_(tinst),
20 4128 : tinfo_(tinfo) {}
21 :
22 :
23 : template<typename InputT, typename ResultT, typename ExternalT, typename SubResultT>
24 : void
25 2048 : WorkStage<InputT, ResultT, ExternalT, SubResultT>::Start(uint32_t stage, FinFn finFn,
26 : const boost::shared_ptr<InputT> & inp) {
27 2048 : assert(!running_);
28 2048 : assert(!finished_);
29 2048 : stage_ = stage;
30 2048 : inp_ = inp;
31 2064 : finFn_ = finFn;
32 2064 : remainingInst_ = tinfo_.size();
33 111821 : for (uint32_t tk = 0 ; tk < tinfo_.size(); tk++) {
34 219360 : workers_.push_back(boost::shared_ptr<WorkProcessor<InputT,SubResultT,ExternalT> >(
35 219582 : new WorkProcessor<InputT,SubResultT,ExternalT>(stage_, efn_,
36 : boost::bind(&WorkStage<InputT,ResultT,ExternalT,SubResultT>::WorkProcCb,
37 : this, tk, _1),
38 109878 : *inp_, tk, tinfo_[tk].first, tinfo_[tk].second)));
39 : }
40 112033 : for (uint32_t tk = 0 ; tk < tinfo_.size(); tk++) {
41 109974 : workers_[tk]->Start();
42 : }
43 :
44 2094 : }
45 :
46 : template<typename InputT, typename ResultT, typename ExternalT, typename SubResultT>
47 : boost::shared_ptr<ResultT>
48 4000 : WorkStage<InputT, ResultT, ExternalT, SubResultT>::Result() const {
49 4000 : if (!finished_) return boost::shared_ptr<ResultT>();
50 4000 : return res_;
51 : }
52 :
53 : template<typename InputT, typename ResultT, typename ExternalT, typename SubResultT>
54 : void
55 2064 : WorkStage<InputT, ResultT, ExternalT, SubResultT>::Release() {
56 2064 : assert(finished_);
57 2064 : assert(subRes_.size() == tinfo_.size());
58 2064 : inp_.reset();
59 111744 : for (uint32_t i = 0; i < subRes_.size(); i++) {
60 109645 : subRes_[i].reset();
61 : }
62 2057 : res_.reset();
63 2064 : }
64 :
65 :
66 : template<typename InputT, typename ResultT, typename ExternalT, typename SubResultT>
67 : bool
68 1936 : WorkStage<InputT, ResultT, ExternalT, SubResultT>::Runner(void) {
69 1936 : running_ = true;
70 1936 : if (!(merger_)(subRes_, inp_, *res_)) {
71 0 : return false;
72 : }
73 1937 : running_ = false;
74 1937 : finished_ = true;
75 1937 : finFn_(true);
76 1935 : return true;
77 : }
78 :
79 :
80 : template<typename InputT, typename ResultT, typename ExternalT, typename SubResultT>
81 : void
82 108518 : WorkStage<InputT, ResultT, ExternalT, SubResultT>::WorkProcCb(uint32_t inst, bool ret_code) {
83 108518 : uint32_t prev = remainingInst_.fetch_sub(1);
84 108518 : if (prev == 1) {
85 2064 : assert(workers_.size() == tinfo_.size());
86 111802 : for (uint32_t i = 0; i < workers_.size(); i++) {
87 109663 : subRes_.push_back(workers_[i]->Result());
88 109367 : workers_[i]->Release();
89 : }
90 2058 : StageProceed(boost::is_same<ResultT,SubResultT>());
91 2064 : if (merger_.empty()) {
92 127 : finished_ = true;
93 127 : finFn_(true);
94 : } else {
95 1937 : res_.reset(new ResultT);
96 1937 : PipelineWorker *w = new PipelineWorker(tid_, tinst_,
97 : boost::bind(&WorkStage<InputT,ResultT,ExternalT,SubResultT>::Runner,
98 : this));
99 1936 : TaskScheduler *scheduler = TaskScheduler::GetInstance();
100 1936 : scheduler->Enqueue(w);
101 : }
102 : }
103 108518 : }
104 :
105 :
106 : template<typename T0, typename T1, typename T2, typename T3,
107 : typename T4, typename T5, typename T6>
108 1937 : WorkPipeline<T0,T1,T2,T3,T4,T5,T6>::WorkPipeline(
109 : WorkStageIf<T0,T1> * s0,
110 : WorkStageIf<T1,T2> * s1,
111 : WorkStageIf<T2,T3> * s2,
112 : WorkStageIf<T3,T4> * s3,
113 : WorkStageIf<T4,T5> * s4,
114 : WorkStageIf<T5,T6> * s5) :
115 1937 : finished_(false),
116 1937 : sg_(
117 3874 : boost::shared_ptr<WorkStageIf<T0,T1> >(s0),
118 3874 : boost::shared_ptr<WorkStageIf<T1,T2> >(s1),
119 3874 : boost::shared_ptr<WorkStageIf<T2,T3> >(s2),
120 3874 : boost::shared_ptr<WorkStageIf<T3,T4> >(s3),
121 3874 : boost::shared_ptr<WorkStageIf<T4,T5> >(s4),
122 3874 : boost::shared_ptr<WorkStageIf<T5,T6> >(s5)) {}
123 :
124 : template<typename T0, typename T1, typename T2, typename T3,
125 : typename T4, typename T5, typename T6>
126 : void
127 1937 : WorkPipeline<T0,T1,T2,T3,T4,T5,T6>::Start(FinFn finFn, const boost::shared_ptr<T0> & inp) {
128 1937 : inp_ = inp;
129 1937 : finFn_ = finFn;
130 1937 : typename boost::tuples::element<0, sg_type>::type g(boost::get<0>(sg_));
131 3874 : g->Start(0, boost::bind(&SelfT::WorkStageCb,
132 1937 : this, 0, _1),inp_);
133 1937 : }
134 :
135 : template<typename T0, typename T1, typename T2, typename T3,
136 : typename T4, typename T5, typename T6>
137 : boost::shared_ptr<typename WorkPipeline<T0,T1,T2,T3,T4,T5,T6>::ResT>
138 1936 : WorkPipeline<T0,T1,T2,T3,T4,T5,T6>::Result() const {
139 1936 : if (!finished_) return boost::shared_ptr<ResT>();
140 1936 : return res_;
141 : }
142 :
143 : template<typename T0, typename T1, typename T2, typename T3,
144 : typename T4, typename T5, typename T6>
145 : void
146 2064 : WorkPipeline<T0,T1,T2,T3,T4,T5,T6>::WorkStageCb(uint32_t stage, bool ret_code) {
147 2064 : switch(stage) {
148 1937 : case 0: {
149 1937 : NextStage<0,T1>();
150 : }
151 1937 : break;
152 127 : case 1: {
153 127 : NextStage<1,T2>();
154 : }
155 127 : break;
156 0 : case 2: {
157 0 : NextStage<2,T3>();
158 : }
159 0 : break;
160 0 : case 3: {
161 0 : NextStage<3,T4>();
162 : }
163 0 : break;
164 0 : case 4: {
165 0 : NextStage<4,T5>();
166 : }
167 0 : break;
168 0 : case 5: {
169 0 : typename boost::tuples::element<5, sg_type>::type g(boost::get<5>(sg_));
170 0 : res_ = g->Result();
171 0 : g->Release();
172 0 : finished_ = true;
173 0 : finFn_(true);
174 0 : }
175 0 : break;
176 : }
177 2064 : }
178 :
179 : template<typename T0, typename T1, typename T2, typename T3,
180 : typename T4, typename T5, typename T6>
181 : template<int kS,typename NextT>
182 : void
183 2064 : WorkPipeline<T0,T1,T2,T3,T4,T5,T6>::NextStage() {
184 2064 : boost::shared_ptr<NextT> res = boost::get<kS>(sg_)->Result();
185 2064 : PipeProceed<kS,boost::is_same<NextT,ResT>::value>::Do(this);
186 2064 : boost::get<kS>(sg_)->Release();
187 2064 : if (boost::get<kS+1>(sg_)) {
188 127 : res_.reset();
189 127 : typename boost::tuples::element<kS+1, sg_type>::type g(boost::get<kS+1>(sg_));
190 127 : g->Start(kS+1, boost::bind(&SelfT::WorkStageCb,
191 : this, kS+1, _1),res);
192 127 : } else {
193 1936 : finished_ = true;
194 1936 : finFn_(true);
195 : }
196 2064 : }
197 :
198 : #endif
|