fec_manager.h 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496
  1. /*
  2. * fec_manager.h
  3. *
  4. * Created on: Sep 27, 2017
  5. * Author: root
  6. */
  7. #ifndef FEC_MANAGER_H_
  8. #define FEC_MANAGER_H_
  9. #include "common.h"
  10. #include "log.h"
  11. #include "lib/rs.h"
  12. const int max_blob_packet_num=30000;//how many packet can be contain in a blob_t ,can be set very large
  13. const u32_t anti_replay_buff_size=30000;//can be set very large
  14. const int max_fec_packet_num=255;// this is the limitation of the rs lib
  15. extern u32_t fec_buff_num;
  16. const int rs_str_len=max_fec_packet_num*10+100;
  17. extern int header_overhead;
  18. extern int debug_fec_enc;
  19. extern int debug_fec_dec;
  20. struct fec_parameter_t
  21. {
  22. int version=0;
  23. int mtu=default_mtu;
  24. int queue_len=200;
  25. int timeout=8*1000;
  26. int mode=0;
  27. int rs_cnt=0;
  28. struct rs_parameter_t //parameters for reed solomon
  29. {
  30. unsigned char x;//AKA fec_data_num (x should be same as <index of rs_par>+1 at the moment)
  31. unsigned char y;//fec_redundant_num
  32. }rs_par[max_fec_packet_num+10];
  33. int rs_from_str(char * s)//todo inefficient
  34. {
  35. vector<string> str_vec=string_to_vec(s,",");
  36. if(str_vec.size()<1)
  37. {
  38. mylog(log_warn,"failed to parse [%s]\n",s);
  39. return -1;
  40. }
  41. vector<rs_parameter_t> par_vec;
  42. for(int i=0;i<(int)str_vec.size();i++)
  43. {
  44. rs_parameter_t tmp_par;
  45. string &tmp_str=str_vec[i];
  46. int x,y;
  47. if(sscanf((char *)tmp_str.c_str(),"%d:%d",&x,&y)!=2)
  48. {
  49. mylog(log_warn,"failed to parse [%s]\n",tmp_str.c_str());
  50. return -1;
  51. }
  52. if(x<1||y<0||x+y>max_fec_packet_num)
  53. {
  54. mylog(log_warn,"invaild value x=%d y=%d, x should >=1, y should >=0, x +y should <%d\n",x,y,max_fec_packet_num);
  55. return -1;
  56. }
  57. tmp_par.x=x;
  58. tmp_par.y=y;
  59. par_vec.push_back(tmp_par);
  60. }
  61. assert(par_vec.size()==str_vec.size());
  62. int found_problem=0;
  63. for(int i=1;i<(int)par_vec.size();i++)
  64. {
  65. if(par_vec[i].x<=par_vec[i-1].x)
  66. {
  67. mylog(log_warn,"error in [%s], x in x:y should be in ascend order\n",s);
  68. return -1;
  69. }
  70. int now_x=par_vec[i].x;
  71. int now_y=par_vec[i].y;
  72. int pre_x=par_vec[i-1].x;
  73. int pre_y=par_vec[i-1].y;
  74. double now_ratio=double(par_vec[i].y)/par_vec[i].x;
  75. double pre_ratio=double(par_vec[i-1].y)/par_vec[i-1].x;
  76. if(pre_ratio+0.0001<now_ratio)
  77. {
  78. if(found_problem==0)
  79. {
  80. mylog(log_warn,"possible problems: %d/%d<%d/%d",pre_y,pre_x,now_y,now_x);
  81. found_problem=1;
  82. }
  83. else
  84. {
  85. log_bare(log_warn,", %d/%d<%d/%d",pre_y,pre_x,now_y,now_x);
  86. }
  87. }
  88. }
  89. if(found_problem)
  90. {
  91. log_bare(log_warn," in %s\n",s);
  92. }
  93. { //special treatment for first parameter
  94. int x=par_vec[0].x;
  95. int y=par_vec[0].y;
  96. for(int i=1;i<=x;i++)
  97. {
  98. rs_par[i-1].x=i;
  99. rs_par[i-1].y=y;
  100. }
  101. }
  102. for(int i=1;i<(int)par_vec.size();i++)
  103. {
  104. int now_x=par_vec[i].x;
  105. int now_y=par_vec[i].y;
  106. int pre_x=par_vec[i-1].x;
  107. int pre_y=par_vec[i-1].y;
  108. rs_par[now_x-1].x=now_x;
  109. rs_par[now_x-1].y=now_y;
  110. double now_ratio=double(par_vec[i].y)/par_vec[i].x;
  111. double pre_ratio=double(par_vec[i-1].y)/par_vec[i-1].x;
  112. //double k= double(now_y-pre_y)/double(now_x-pre_x);
  113. for(int j=pre_x+1;j<=now_x-1;j++)
  114. {
  115. int in_x=j;
  116. //////// int in_y= double(pre_y) + double(in_x-pre_x)*k+ 0.9999;// round to upper
  117. double distance=now_x-pre_x;
  118. /////// double in_ratio=pre_ratio*(1.0-(in_x-pre_x)/distance) + now_ratio *(1.0- (now_x-in_x)/distance);
  119. ////// int in_y= in_x*in_ratio + 0.9999;
  120. int in_y= pre_y +(now_y-pre_y) *(in_x-pre_x)/distance +0.9999;
  121. if(in_x+in_y>max_fec_packet_num)
  122. {
  123. in_y=max_fec_packet_num-in_x;
  124. assert(in_y>=0&&in_y<=max_fec_packet_num);
  125. }
  126. rs_par[in_x-1].x=in_x;
  127. rs_par[in_x-1].y=in_y;
  128. }
  129. }
  130. rs_cnt=par_vec[par_vec.size()-1].x;
  131. return 0;
  132. }
  133. char *rs_to_str()//todo inefficient
  134. {
  135. static char res[rs_str_len];
  136. string tmp_string;
  137. char tmp_buf[100];
  138. assert(rs_cnt>=1);
  139. for(int i=0;i<rs_cnt;i++)
  140. {
  141. sprintf(tmp_buf,"%d:%d",int(rs_par[i].x),int(rs_par[i].y));
  142. if(i!=0)
  143. tmp_string+=",";
  144. tmp_string+=tmp_buf;
  145. }
  146. strcpy(res,tmp_string.c_str());
  147. return res;
  148. }
  149. rs_parameter_t get_tail()
  150. {
  151. assert(rs_cnt>=1);
  152. return rs_par[rs_cnt-1];
  153. }
  154. int clone(fec_parameter_t & other)
  155. {
  156. version=other.version;
  157. mtu=other.mtu;
  158. queue_len=other.queue_len;
  159. timeout=other.timeout;
  160. mode=other.mode;
  161. assert(other.rs_cnt>=1);
  162. rs_cnt=other.rs_cnt;
  163. memcpy(rs_par,other.rs_par,sizeof(rs_parameter_t)*rs_cnt);
  164. return 0;
  165. }
  166. int clone_fec(fec_parameter_t & other)
  167. {
  168. assert(other.rs_cnt>=1);
  169. rs_cnt=other.rs_cnt;
  170. memcpy(rs_par,other.rs_par,sizeof(rs_parameter_t)*rs_cnt);
  171. version++;
  172. return 0;
  173. }
  174. };
  175. extern fec_parameter_t g_fec_par;
  176. //extern int dynamic_update_fec;
  177. const int anti_replay_timeout=120*1000;// 120s
  178. struct anti_replay_t
  179. {
  180. struct info_t
  181. {
  182. my_time_t my_time;
  183. int index;
  184. };
  185. u64_t replay_buffer[anti_replay_buff_size];
  186. unordered_map<u32_t,info_t> mp;
  187. int index;
  188. anti_replay_t()
  189. {
  190. clear();
  191. }
  192. int clear()
  193. {
  194. memset(replay_buffer,-1,sizeof(replay_buffer));
  195. mp.clear();
  196. mp.rehash(anti_replay_buff_size*3);
  197. index=0;
  198. return 0;
  199. }
  200. void set_invaild(u32_t seq)
  201. {
  202. if(is_vaild(seq)==0)
  203. {
  204. mylog(log_trace,"seq %u exist\n",seq);
  205. //assert(mp.find(seq)!=mp.end());
  206. //mp[seq].my_time=get_current_time_rough();
  207. return;
  208. }
  209. if(replay_buffer[index]!=u64_t(i64_t(-1)))
  210. {
  211. assert(mp.find(replay_buffer[index])!=mp.end());
  212. mp.erase(replay_buffer[index]);
  213. }
  214. replay_buffer[index]=seq;
  215. assert(mp.find(seq)==mp.end());
  216. mp[seq].my_time=get_current_time();
  217. mp[seq].index=index;
  218. index++;
  219. if(index==int(anti_replay_buff_size)) index=0;
  220. }
  221. int is_vaild(u32_t seq)
  222. {
  223. if(mp.find(seq)==mp.end()) return 1;
  224. if(get_current_time()-mp[seq].my_time>anti_replay_timeout)
  225. {
  226. replay_buffer[mp[seq].index]=u64_t(i64_t(-1));
  227. mp.erase(seq);
  228. return 1;
  229. }
  230. return 0;
  231. }
  232. };
  233. struct blob_encode_t
  234. {
  235. char input_buf[(max_fec_packet_num+5)*buf_len];
  236. int current_len;
  237. int counter;
  238. char *output_buf[max_fec_packet_num+100];
  239. blob_encode_t();
  240. int clear();
  241. int get_num();
  242. int get_shard_len(int n);
  243. int get_shard_len(int n,int next_packet_len);
  244. int input(char *s,int len); //len=use len=0 for second and following packet
  245. int output(int n,char ** &s_arr,int & len);
  246. };
  247. struct blob_decode_t
  248. {
  249. char input_buf[(max_fec_packet_num+5)*buf_len];
  250. int current_len;
  251. int last_len;
  252. int counter;
  253. char *output_buf[max_blob_packet_num+100];
  254. int output_len[max_blob_packet_num+100];
  255. blob_decode_t();
  256. int clear();
  257. int input(char *input,int len);
  258. int output(int &n,char ** &output,int *&len_arr);
  259. };
  260. class fec_encode_manager_t:not_copy_able_t
  261. {
  262. private:
  263. u32_t seq;
  264. //int fec_mode;
  265. //int fec_data_num,fec_redundant_num;
  266. //int fec_mtu;
  267. //int fec_queue_len;
  268. //int fec_timeout;
  269. fec_parameter_t fec_par;
  270. my_time_t first_packet_time;
  271. my_time_t first_packet_time_for_output;
  272. blob_encode_t blob_encode;
  273. char input_buf[max_fec_packet_num+5][buf_len];
  274. int input_len[max_fec_packet_num+100];
  275. char *output_buf[max_fec_packet_num+100];
  276. int output_len[max_fec_packet_num+100];
  277. int counter;
  278. //int timer_fd;
  279. //u64_t timer_fd64;
  280. int ready_for_output;
  281. u32_t output_n;
  282. int append(char *s,int len);
  283. ev_timer timer;
  284. struct ev_loop *loop=0;
  285. void (*cb) (struct ev_loop *loop, struct ev_timer *watcher, int revents)=0;
  286. public:
  287. fec_encode_manager_t();
  288. ~fec_encode_manager_t();
  289. fec_parameter_t & get_fec_par()
  290. {
  291. return fec_par;
  292. }
  293. void set_data(void * data)
  294. {
  295. timer.data=data;
  296. }
  297. void set_loop_and_cb(struct ev_loop *loop,void (*cb) (struct ev_loop *loop, struct ev_timer *watcher, int revents))
  298. {
  299. this->loop=loop;
  300. this->cb=cb;
  301. ev_init(&timer,cb);
  302. }
  303. int clear_data()
  304. {
  305. counter=0;
  306. blob_encode.clear();
  307. ready_for_output=0;
  308. seq=(u32_t)get_fake_random_number(); //TODO temp solution for a bug.
  309. if(loop)
  310. {
  311. ev_timer_stop(loop,&timer);
  312. }
  313. return 0;
  314. }
  315. int clear_all()
  316. {
  317. //itimerspec zero_its;
  318. //memset(&zero_its, 0, sizeof(zero_its));
  319. //timerfd_settime(timer_fd, TFD_TIMER_ABSTIME, &zero_its, 0);
  320. if(loop)
  321. {
  322. ev_timer_stop(loop,&timer);
  323. loop=0;
  324. cb=0;
  325. }
  326. clear_data();
  327. return 0;
  328. }
  329. my_time_t get_first_packet_time()
  330. {
  331. return first_packet_time_for_output;
  332. }
  333. int get_pending_time()
  334. {
  335. return fec_par.timeout;
  336. }
  337. int get_type()
  338. {
  339. return fec_par.mode;
  340. }
  341. //u64_t get_timer_fd64();
  342. int reset_fec_parameter(int data_num,int redundant_num,int mtu,int pending_num,int pending_time,int type);
  343. int input(char *s,int len/*,int &is_first_packet*/);
  344. int output(int &n,char ** &s_arr,int *&len);
  345. };
  346. struct fec_data_t
  347. {
  348. int used;
  349. u32_t seq;
  350. int type;
  351. int data_num;
  352. int redundant_num;
  353. int idx;
  354. char buf[buf_len];
  355. int len;
  356. };
  357. struct fec_group_t
  358. {
  359. int type=-1;
  360. int data_num=-1;
  361. int redundant_num=-1;
  362. int len=-1;
  363. int fec_done=0;
  364. //int data_counter=0;
  365. map<int,int> group_mp;
  366. };
  367. class fec_decode_manager_t:not_copy_able_t
  368. {
  369. anti_replay_t anti_replay;
  370. fec_data_t *fec_data=0;
  371. unordered_map<u32_t, fec_group_t> mp;
  372. blob_decode_t blob_decode;
  373. int index;
  374. int output_n;
  375. char ** output_s_arr;
  376. int * output_len_arr;
  377. int ready_for_output;
  378. char *output_s_arr_buf[max_fec_packet_num+100];//only for type=1,for type=0 the buf inside blot_t is used
  379. int output_len_arr_buf[max_fec_packet_num+100];//same
  380. public:
  381. fec_decode_manager_t()
  382. {
  383. fec_data=new fec_data_t[fec_buff_num+5];
  384. assert(fec_data!=0);
  385. clear();
  386. }
  387. /*
  388. fec_decode_manager_t(const fec_decode_manager_t &b)
  389. {
  390. assert(0==1);//not allowed to copy
  391. }*/
  392. ~fec_decode_manager_t()
  393. {
  394. mylog(log_debug,"fec_decode_manager destroyed\n");
  395. if(fec_data!=0)
  396. {
  397. mylog(log_debug,"fec_data freed\n");
  398. delete fec_data;
  399. }
  400. }
  401. int clear()
  402. {
  403. anti_replay.clear();
  404. mp.clear();
  405. mp.rehash(fec_buff_num*3);
  406. for(int i=0;i<(int)fec_buff_num;i++)
  407. fec_data[i].used=0;
  408. ready_for_output=0;
  409. index=0;
  410. return 0;
  411. }
  412. //int re_init();
  413. int input(char *s,int len);
  414. int output(int &n,char ** &s_arr,int* &len_arr);
  415. };
  416. #endif /* FEC_MANAGER_H_ */