-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathxfiber.cpp
More file actions
1261 lines (1000 loc) · 52.3 KB
/
Copy pathxfiber.cpp
File metadata and controls
1261 lines (1000 loc) · 52.3 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
#include <stdio.h>
#include <errno.h> //是C语言C标准函式库里的标头档,定义了通过错误码来回报错误资讯的宏。
#include <error.h>//error系列函数是Linux系统编程中,一种debug的方式,
#include <cstring>
#include <iostream>
#include <sys/epoll.h>//epoll
#include <sys/types.h>
#include <signal.h>
#include <unistd.h>
#include <assert.h>//断言
#include "xfiber.h"
XFiber::XFiber() {
//XFiber 调度器初始化时 没有任何fiber被调度
this->cur_dispatch_fiber_ = nullptr;
//>注意<: 这里是deque
this->ready_fibers_ = std::deque<Fiber *>();
this->running_fibers_ = std::deque<Fiber *>();
// epoll 初始化
//>注意<: epoll_create 和 epoll_create1
this->ep_fd_ = epoll_create1(0);//
if (ep_fd_ < 0) {
LOG_ERROR("epoll_create failed, msg=%s", strerror(errno));
exit(-1);
}
/*
睁大眼睛: epoll_create和epoll_create1
第一级调用: int epoll_create(int size); size参数 只是对内核的建议,没啥用
第二级调用: int epoll_create1(int flags);
第三级调用: ep_alloc()创建内部数据(eventpoll)
-------------------------------
+ 1.初始化epoll文件等待队列(双向链表)
+ 2.初始化eventpoll文件唤醒队列(双向链表)
+ 3.初始化就绪队列(双向链表)
-----------------------------------
拓展: epoll_create1的参数flag
flag == 0 : 和 epoll_create() 一样: 创建一个 epoll实例
flag == EPOLL_CLOEXEC
在文件描述符上面设置执行时关闭(FD_CLOEXEC)标志描述符。
*/
}
XFiber::~XFiber(){
//关掉epoll
close(this->ep_fd_);
}
//获取正在 xfiber上调度的fiber的ctx (fiber上了xfiber后 ctx替换)
//>注意<: fiber 上了xfiber调度后 xfiber 的sched_ctx_ 存储的 是fiber的上下文信息
XFiberCtx* XFiber::SchedCtx(){
return &sched_ctx_;
}
//>注意<: 新建一个fiber, 并且只是简单的将他加入到ready队列中去
//>注意<: std::function<void ()> run 是为了适应仿函数
void XFiber::CreateFiber(std::function<void ()> run, size_t stack_size, std::string fiber_name) {
if (stack_size == 0) {
stack_size = 1024 * 1024;
}
Fiber *fiber = new Fiber(run, this, stack_size, fiber_name);
ready_fibers_.push_back(fiber);
LOG_DEBUG("create a new fiber with id[%lu]", fiber->Seq());
}
// 唤醒某一个协程,上调度器
// 唤醒时,需要处理: fiber协程对应的 可读可写事件 以及 超时事件
// 如何处理: 1. 该fiber加入就绪队列
// 2. 从 xfiber 的 io_waitiing_events中 删除 该fiber 负责的 可读可写事件 的fd 的映射关系
// |==> 为什么要这么做: 一个套接字fd 被epoll侦听,该fd 对应的fiber 只有一个
// |==> ==> 满足了 一个fd 只能被一个fiber处理的 情况
// |==> 但是 ! 一个fiber 也可能同时 负责处理很多个fd
// |==> ==> 所以 当该fiber 被任何一个epoll侦听到的 fd 唤醒后, 需要先 将本身负责的所有的读写fd ,在io_waiting_events 映射中删除
// |==> ==> ==> 一个fd 只能被一个fiber处理,一个fiber可以处理多个fd. 我一旦被触发,多余的重复映射就删除吊
// |==> ==> ==> 我先起床,我先把我能吃的早餐 全部吃掉
// 3. 从 xfiber 的 超时队列 中 删除 该fiber 对应的 超时事件
// 4. 从epoll中踢出: 防止二次唤醒
// =======> 理由 wakeup唤醒后 加入 dispatch ,按顺序上调度器, 上了调度器之后 会 处理fiber自身的读写超时事件
// ===> 处理完 ,本来就应该删除,只不过这里是提前了一点点而已
void XFiber::WakeUpFiber(Fiber * fiber) {
// DEBUG_ENABLE : 默认为0 LOG_DEBUG(fmt,...)
// LOG_DEBUG("try wakeup fiber[%lu] %p", fiber->Seq(), fiber);
//第一步: 被唤醒的 fiber ,要先加入就绪队列
this->ready_fibers_.push_back(fiber);
LOG_INFO("try wakeup fiber[%lu] %p, %d fibers ready to run", fiber->Seq(), fiber, ready_fibers_.size());
//第二步: 从io_waiting_events中删除该fiber负责 的所有的 多个的 读写fd 和其他fiber的映射关系
//>提问<: 等待队列是哪个? running是运行 ,ready是就绪
//>解答<: 等待队列 是 epoll 监听fd 和对应fiber对象的map 即io_waiting_events
// io_waiting_events: 里面存储的是所有的 fd和fiber的映射关系
// epoll触发fd之后,从该map中寻找对应的 fiber
// 2.1 先获取 fiber的等待事件集合 ==> waiting_events
// >重点<: 一个 fd一般只能交给 一个 协程(fiber)处理
// +_+ 但是 一个 fiber 可以处理多个fd ==> 所以 我们需要去除多余的 重复的映射关系
//>重点<: 一个fiber 可以对应很多个 很多个 fd事件 每个fd事件交给一个 新的 fiber去处理
// ==> 所以这样就有了 waitingevents
//每一个fiber 都有一个 waitingevents 里面包含了 这个fiber 负责的fd对象
// 现在 这个fiber即将被 wakeup 所以也需要 将这些fd处理掉
auto & waiting_events = fiber->GetWaitingEvents();
// 2.2 在 对 waiting_events 内 所有的 fd 进行处理: 可读fd 可写fd 超时队列
// events:事件集合包含了 读事件 写事件 超时事件, 需要分别删除对应的 映射关系
// 删除 读fd的映射关系
for (size_t i = 0;i < waiting_events.waiting_fds_r_.size();i ++) {
int fd = waiting_events.waiting_fds_r_[i];
//>注意<: io_waiting_fibers_ 就是那个 map 就是全局等待队列
auto iter = io_waiting_fibers_.find(fd);
//如果对应fd 存在 等待队列中, 就把该fd 删除
if (iter != io_waiting_fibers_.end()) {
io_waiting_fibers_.erase(iter);
}
}
// 处理写事件的映射关系
for (size_t i = 0;i < waiting_events.waiting_fds_w_.size();i ++) {
int fd = waiting_events.waiting_fds_r_[i];
auto iter = io_waiting_fibers_.find(fd);
if (iter != io_waiting_fibers_.end()) {
io_waiting_fibers_.erase(iter);
}
}
//第三步: 从 全局超时队列expire_events_(超时时间,fiber)组成的map 中 删除掉对应的fd事件
//获取 fiber的事件集合设置的超时时间 ==> 一旦到达该时间 整个fiber超时作废
//>提问<: 为什么要在这里处理超时问题?
//>解答<: 超时机制本质: 时间戳 和当前时间做对比: 超时:则整个超时队列会wakeup
//>重点<: 部分超时的 会走自身内部 超时逻辑, 未超时的 相当于提前唤醒
//>重点<: 一个fiber 可能会负责 多个fd, 超时属性是fd 和fiber的, 一旦当前fd 的负责人和本fiber 匹配,则fd的超时时间,也就是和本fiber匹配,则去掉本fiber的超市属性
WaitingEvents &evs = fiber->GetWaitingEvents(); // 获取读事件fd,检测 是否设置有超时时间 , 然后去超时队列中, 寻找,当前fiber 是否有超时属性
/* 二者合二为一
for (size_t i = 0; i < evs.waiting_fds_r_.size(); i++) {
auto expired_at = evs.expire_at_;
if (expired_at > 0) {
auto expired_iter = expire_events_.find(expired_at);
//如果 该fiber 没有超时属性
if (expired_iter->second.find(fiber) == expired_iter->second.end()) {
LOG_ERROR("not fiber [%lu] in expired events", fiber->Seq());
}
else {
// 如果该fiber 有超时属性, 则 因为被提前唤醒,而失去超时属性 ===> 解除超时映射
expired_iter->second.erase(fiber);
}
}
}
for (size_t i = 0; i < evs.waiting_fds_w_.size(); i++) {
int64_t expired_at = evs.expire_at_;
if (expired_at > 0) {
auto expired_iter = expire_events_.find(expired_at);
if (expired_iter->second.find(fiber) == expired_iter->second.end()) {
LOG_ERROR("not fiber [%lu] in expired events", fiber->Seq());
}
else {
expired_iter->second.erase(fiber);
}
}
}
*/
//>重点<: 从这里可以看出 超时时间 是 fiber的属性,
//>注意<: 这里 比较的是当前fiber的超时时间,和超时队列里 其余fiber的超时时间
int64_t expire_at = waiting_events.expire_at_;
if (expire_at > 0) {
// 这一步: 有没有和我具有相同超时时间的其他fiber对象呢?
auto expired_iter = expire_events_.find(expire_at);
if (expired_iter->second.find(fiber) == expired_iter->second.end()) {
LOG_WARNING("not fiber [%lu] in expired events", fiber->Seq());
}
else {
LOG_DEBUG("remove fiber [%lu] from expire events...", fiber->Seq());
expired_iter->second.erase(fiber);
}
}
// 时间: 2023年3月7日20:22:28 不需要这些
// >重点<: 从 epoll 中也需要踢出
// >解答<: 其实也不用从epoll中踢出,这个要看你的写法吧
// >解答< : 如果不踢出 的话, 就在初始化的时候 ,使用epoll_ctl_mod
// 因为 epoll_ctl_add 如果epoll内部该fd被触发后 fd还存在,但是事件集合被清空了,所以add会失败
// epoll 照着set ,把里面的fd踢出
// >提问<: 为什么要从epoll中踢出啊?
/* >解答<: epoll踢出:
>重点<: epoll_wait工作原理: 这里的epoll_ctl 并没有使用epollshort 所以不需要重新注册
等侍注册在epfd上的socket fd的事件的发生,
如果发生则将 `发生的sokct fd`和 `事 件类型放入到events数组中`
`并 且将注册在epfd上的socket fd的事件类型给清空`,
所以如果下一个循环你还要关注这个socket fd的话,则需要用epoll_ctl(epfd,EPOLL_CTL_MOD,listenfd,&ev)来重新设置socket fd的事件类型。
`这时不用EPOLL_CTL_ADD,因为socket fd并未清空`,只是 `事件类型清空`。这一步非常重要。
所以: wakeup之后 et 模式 + 死循环, 保证了 事件平稳的处理完成
处理完成之后,fd仍然存在,但是 fd的事件 会被清空,
>重点<: ====> 下一次相同fd 初始化时 epoll_ctl_add 会失败
====> 如果没有相同的 fd继续来, 那么epoll中的 fd也太多了 ==> 红黑树虽然没啥影响
*/
/*
for (auto iter : waiting_fds) {
if (epoll_ctl(this->ep_fd_, EPOLL_CTL_DEL, iter, nullptr) < 0) {
perror("epoll_ctl");
}
}
*/
// 4. 删除信号量
for (auto sem : waiting_events.waiting_sems_) {
auto iter = sem_infos_.find(sem);
if (iter != sem_infos_.end()) {
iter->second.fibers_.erase(fiber);
}
}
LOG_DEBUG("fiber [%lu] %p has wakeup success, ready to run!", fiber->Seq(), fiber);
}
//>重点<: xfiber调度器 调度策略
void XFiber::Dispatch() {
//>注意<: 死循环! 也就是说 上下文一直都在这个循环内,无法跳出Dispatch
// ==> fiber上调度器 后也会一直在 这个 dispatch 内运行
while(true) {
// 情况1: 就绪队列 不为空 ==> 处理就绪队列
if (this->ready_fibers_.size() > 0 ) {
// 完美转发: 实现==>万能引用 + 应用折叠
// 功能: 将所有的 准备队列 全部 转换成就绪队列, 从而运行
this->running_fibers_ = std::move(this->ready_fibers_);
this->ready_fibers_.clear();
LOG_DEBUG("there are %ld fiber(s) in ready list, ready to run...", running_fibers_.size());
//>注意<: 现在需要处理所有的 running队列
// 策略: 先来先服务
// 遍历运行队列 ,处理每一个fiber
for (auto iter = this->running_fibers_.begin();iter != this->running_fibers_.end();iter ++) {
Fiber * fiber = * iter; // 获取即将上调度器的fiber
this->cur_dispatch_fiber_ = fiber; //>注意<: 提前保存该对象,考虑到调度途中 万一需要让出cpu(暨YILED)
LOG_DEBUG("switch from sched to fiber[%lu]", fiber->Seq());
// assert(SwitchCtx(SchedCtx(), fiber->Ctx()) == 0);
//#define SwitchCtx(from, to) swapcontext(from, to) 交换ctx
//>重点<: XFiber 调度器 如何调度? swapcontext来进行CPU上下文切换
/*
解释:
1.当前cpu上运行的是XFIBER调度器,
==> 所以当前cpu上下文也是 XFIBER的上下文
2.XFIBER即将调度fiber对象上cpu运行,
==> 所以就需要保存 当前 XFIBER的 上下文信息, swapcontext(当前XFIBER的ctx,fiber的ctx)
==> 交换后: xfiber的ctx 就是 fiber的ctx ,暨相当于调度上了fiber
==> cur_distance_fiber 就是 fiber的ctx
==> 并且 使用 cur_distance_fiber 的上下文信息来进行运行 fiber对象
3. 由于 交换后的 cur_dispacth_fib_.fib_ctx.uc_link == 最开始的xfiber.xfiber_ctx
==> 暨: 交换后的 cur_dispatch_fib 运行完后 会上行 最开始的xfiber.xfiber_ctx , 暨继续运行 xfiber调度器上次交换出去的地方
==> 所以 在 FIBER对象(被调度协程)运行完毕后,会回归到 xfiber(调度协程)的上下文
==> 所以 ,回去之后 的下一条命令 就是 交换后的cur_dispatch_fib_ = nullptr,暨表示 fiber已经运行完毕
*/
assert(SwitchCtx(SchedCtx(),fiber->Get_Fiber_Ctx()) == 0);// 这条语句开始执行是: 是去运行fiber上下文 这条语句执行完毕时: 时从fiber回到xfiber
//注意:
//该调度对象(暨协程)已经彻底啊调度完成
// ===> 为什么说调度完成了?
// 因为他已经 从被调度的fiber对象的上下文中,再一次的返回到了 XFIBER的上下文
// 浏览 FIBER中 ctx的初始化
// fiber->fiber_ctx_.uc_link = xfiber_ctx_
// 任何一个 fiber对象(暨协程),在调度后,都必须回到XFIBER调度器,从而继续调度运行队列里面的后续协程
// 该fiber对象调度完成,就就没啥用了,所以就不需要保存了
this->cur_dispatch_fiber_ = nullptr; //fiber 运行完毕回到xfiber 所以 当前即将运行的fiber 为null
if (fiber->IsFinished()) {
LOG_INFO("fiber[%lu] finished, free it!", fiber->Seq());
delete fiber;//彻底运行完毕 ,就没啥用了
}
}
this->running_fibers_.clear();//运行队列中所有的被调度FIBER(暨被调度协程)已经调度完成
}
// 情况2: 在全局超时队列中处理带有超时时间的事件
int64_t now_ms = util::NowMs();
//this->expire_events_.begin->first: 暨超时时间(dead_line) 是 超时时间expire_at_ + now_ms(设置时的事件戳) 计算而来
// while (!this->expire_events_.empty() && this->expire_events_.begin()->first <= now_ms) {
//>注意<: 这里特别 离谱!
//>重点<: 这里超级离谱!
//>提问<: 这里为什么是这样写? 为什么不能写 >=?
//>解答<: 部分超时 and 全部超时
// --> begin <= 当前时间 部分超时: 部分超时后,处理整个超时队列
// --> 超时的fiber 进行唤醒,上cpu 执行内部自己的超时逻辑
// --> 未超时的fiber: 就相当于 是提前唤醒
// --> begin >= 当前时间 全部都没有超时: 则彻底忽略整个 超时队列
//>解答<: 超时队列 内部 未超时的 fiber 会被提前唤醒,失去超时属性, 一旦全部超时,则根本不进行处理.
// ================> 不进行处理的 后果:
// ================> 举例: 一个fiber 负责的fd 有超时时间, 但是已经超时很久了, 一直在超市队列里面等着
// ================> 直到 超时队列.begin <= 当前时间 ==> 视为 超时 ,然后 被唤醒 去执行自己的超时逻辑 之后推出
//>重点<: :1. 如果最小的 超时时限都没有 超时, 那么全部都没有超时
// 2. 如果 反写乘 >= 会导致所有未超时的fiber都会被再次上cpu,哪怕根本没有数据流入流出
// 3. 如果 协程 >= 那么 超时时限 根本不会准确: 因为1多次 的上下cpu sleepms = 当前时间 + timeout
// ==> 每一次上cpu timeout都会改变
// 4. 为什么要写成这样? wakeup超时fiber 不是事与愿违吗? 超时的不是应该舍弃吗? wakeup之后 有注定会让 超时fiber上一次CPU? 这不是和初衷(超时舍弃)相违背吗
// ======> 超时的fiber唤醒, 唤醒后wakeup
// ===> wakeup 先加入ready_对列 然后再从超时队列 等待集合 epoll中删除
// ==> 然后再次上cpu ,
//>注意<: 此时: 是最后一次上cpu ,上cpu的目的是为了 //>提示<: 超时fiber进行自己的超时逻辑,退出!,再也找不到他了,彻底finished了
//>重点<: 这样就保证了: 1. 超时队列里面全都是未超时的
// 2. 一旦fiber超时 就无效化(彻底finished了)(哪怕再次被触发也无效)
// 3. 保证了 未超时fiber不会频繁上cpu 导致 超时时间线动态变化 ===>从而影响 到固定的超时时限(死线不会后移了)
// 4. 一旦超时fiber 被最后一次上cpu 会自己去执行自己的内部的超时逻辑
while (!this->expire_events_.empty() && this->expire_events_.begin()->first <= now_ms) {
//std::map<int64_t, std::set<Fiber *>> expire_events_; 按照超时时间排序
// this->expire_events_.begin()->first 红黑树,也就是 最小的那个超时时间
std::set<Fiber *> &expired_fibers = expire_events_.begin()->second;
// 依次处理超时队列
while (!expired_fibers.empty())
{
std::set<Fiber *>::iterator expired_fiber = expired_fibers.begin();
//唤醒超时队列里面的每一个
WakeUpFiber(* expired_fiber);
}
expire_events_.erase(expire_events_.begin());//1个大类的超时时间线fiber被全部唤醒完成
}
// 情况3: 处理epoll监听中被 触发的fd
// ===> 特殊: epoll事件 被触发后 立即调用WakeUpFiber 将该fiber 加入就绪队列 删除 等待集合中的可读可写fiber,
// 删除超时队列里的 fiber
#define MAX_EVENT_COUNT 512
/*结构体细节: 事件+ 存储信息的联合体
记住union的特性,它只会保存最后一个被赋值的成员,
所以不要data.fd,data.ptr都赋值;
通用的做法就是给ptr赋值,
fd只是有些时候为了演示或怎么样罢了
>>>>>>>> 可见,这个 data成员还与具体的使用方式相关。 <<<<<<<<<
typedef union epoll_data
{
void *ptr; : 感兴趣的数据
int fd;
uint32_t u32;
uint64_t u64;
} epoll_data_t;
struct epoll_event
{
// epoll 注册的事件,比如EPOLLIN、EPOLLOUT等等,
//这个参数在epoll_ctl注册事件时,可以明确告知注册事件的类型。
uint32_t events;
epoll_data_t data;//保存触发事件的某个文件描述符相关的数据
} __EPOLL_PACKED;
*/
struct epoll_event evs[MAX_EVENT_COUNT];
/*>注意<: epoll_wait
函数功能: 等待epoll事件的产生
int epoll_wait(int epfd, struct epoll_event * events, int maxevents, int timeout);
参数:
1. epfd: eopll结构句柄
2. epoll_events: 从内核得到的 epoll事件集合,会拷贝一份到该处,注意是地址
3. maxevents: 告之内核这个events有多大
maxevents <= epoll_create()时的size,
4. timeout : 超时时间
① 0: 立即返回
② -1: -1将不确定,也有说法说是永久阻塞
返回值:
1. 成功 : > 0 返回需要处理的事件数目
2. 失败 : -1:
epoll_wait处理完事件后, 会对fd上的 事件进行清空
fd依旧存在,但不再关心这个事件
//。epoll_wait将会接收到消息,并且将数据拷贝到用户空间,清空链表。
对于LT模式epoll_wait清空就绪链表之后会检查该文件描述符是哪一种模式,
如果为LT模式,且必须该节点确实有事件未处理, =============> 继续关心 这个 ,通知多次
那么就会把该节点重新放回到 刚刚删除掉的且刚准备好的就绪链表,
epoll_wait马上返回。ET模式不会检查,只会调用一次 ==========> 不在加入到关心链表,只通知一次
*/
// 设置超时时间2s
int n = epoll_wait(this->ep_fd_, evs, MAX_EVENT_COUNT, 2);
if (n < 0) {
LOG_ERROR("epoll_wait error, msg=%s", strerror(errno));
continue; //整个朱调度器 是一个死循环
}
//处理触发事件
for (int i = 0; i < n; i++) {
//获取事件单个epoll事件
struct epoll_event &ev = evs[i];//evs 是epoll获取到数据 从rdlist 拷贝到此处的
//epoll 事件 文件描述符
int fd = ev.data.fd;
// map中 fd 对应的是 一个waitting_fiber : witing_fiber:中有两份 fiber* : 读fiber对象 写fiber对象
auto fiber_iter = io_waiting_fibers_.find(fd);
//被通知的事件 ,存在于 等待队列中
if (fiber_iter != io_waiting_fibers_.end()) {
// 获取 等待集合中 的等待fiber
WaitingFibers &waiting_fiber = fiber_iter->second;
//可读事件
if (ev.events & EPOLLIN) {
if (waiting_fiber.r_ == nullptr) {
LOG_WARNING("fd[%d] is readfd and this fd is null", fd);
}
else {
LOG_DEBUG("waiting fd[%d] has fired IN event, wake up pending fiber[%lu]", fd, waiting_fiber.r_->Seq());
WakeUpFiber(waiting_fiber.r_); //唤醒fd对应的 读事件处理协程
}
}
else if (ev.events & EPOLLOUT) {
if (waiting_fiber.w_ == nullptr) {
LOG_WARNING("fd[%d] has been fired OUT event, but not found any fiber to handle!", fd);
}
else {
LOG_DEBUG("waiting fd[%d] has fired OUT event, wake up pending fiber[%lu]", fd, waiting_fiber.w_->Seq());
WakeUpFiber(waiting_fiber.w_);
}
}
}
}
}
}
//调用该函数的协程: 主动让出cpu
// 让出cpu的结果就是,会回到XFIBER调度器,然后调度器,去选择下一个调度对象
//>注意<: Yiled 和switchtoscheduler 的区别: 到底是是否加入到就绪队列?
void XFiber::Yield() {
// 当前被调度 的 fiber的 ctx 不为空
assert(this->cur_dispatch_fiber_!= nullptr);// 说明当前调度器上的是fiber协程
LOG_INFO("will Yield , ready_fibers_.size + 1");
// 主动切出的后仍然是ready状态,等待下次调度
ready_fibers_.push_back(this->cur_dispatch_fiber_);
SwitchToScheduler();
}
//函数功能: 从当前的 正在执行的协程(fiber对象)切换出来, 转而去执行xfiber调度器
// 单纯的调用该函数, 不会讲ctx加入到就绪队列 ==> 一旦调用该函数, 就说明 该fiber 已经不再上cpu了
void XFiber::SwitchToScheduler() {
assert(this->cur_dispatch_fiber_ != nullptr);
LOG_DEBUG("switch to sched");
// assert(SwitchCtx(cur_dispatch_fiber_, SchedCtx()) == 0);
// 就是单纯的 swapcontext
// std::cout<< cur_dispatch_fiber_->Get_Fiber_Ctx()<<std::endl;
// std::cout<< SchedCtx()<<std::endl;
assert(SwitchCtx(cur_dispatch_fiber_->Get_Fiber_Ctx(),SchedCtx()) == 0);
}
//设置超时时间
void XFiber::SleepMs(int ms) {
if (ms < 0) {
return;
}
//当前系统时间+ 超时时间 == 死线
int64_t expired_at = util::NowMs() + ms;
// 为什么使用-1? 因为别的fd不可能是负数 为什么是false: sleepms单纯的就是为了 休眠,不参与任何读写
// RegisterFdWithCurrFiber(-1,expired_at,false);
WaitingEvents events;
events.expire_at_ = expired_at;
RegisterWaitingEvents(events);
//切换回调度器
SwitchToScheduler();
}
//>提示<: 第一版 :不能解决信号量的问题
//功能: 将目标fd 和当前fiber相互绑定
// fd 绑定 fiber: <map>io_waiting_ : fd 映射读写fiber
// fiber 绑定 fd: fiber 内部的 读写fd数组 中药新增该fd
// 超时时间: 将 超时时间 和当前fiber的关系 进行绑定 即expire_events
/*
bool XFiber::RegisterFdWithCurrFiber(int fd, int64_t expired_at, bool is_write) {
auto iter = io_waiting_fibers_.find(fd);
// 当前fd 尚未出现在map中
if (iter == io_waiting_fibers_.end()) {
WaitingFibers wb;
if (is_write == false) {
// 说明这是一个 专门写入数据的fiber对象 (暨协程)
// wb是一个keyvalue结构, 保存的 是 fiber对象, fiber对象里面有 ctx上下文,ctx上下文继续去执行
wb.r_ = this->cur_dispatch_fiber_;
wb.w_ = nullptr;
//初始化
this->cur_dispatch_fiber_->SetReadEvent(Fiber::FdEvent(fd));
} else {
wb.w_ = this->cur_dispatch_fiber_;
wb.r_ = nullptr;
this->cur_dispatch_fiber_->SetWriteEvent(Fiber::FdEvent(fd));
}
} else {
// fd已经存在于 map中了 ,该fiber再一次的被调用之后, 上下文信息会发生改变,所以我们需要更新上下文信息
if (is_write == false) {
iter->second.r_ = this->cur_dispatch_fiber_;
this->cur_dispatch_fiber_->SetReadEvent(Fiber::FdEvent(fd));
} else {
iter->second.w_ = this->cur_dispatch_fiber_;
this->cur_dispatch_fiber_->SetWriteEvent(Fiber::FdEvent(fd));
}
}
return true;
}
*/
//>提示<: 第二版 : 解决了信号量注册的问题
void XFiber::RegisterWaitingEvents(WaitingEvents &events) {
assert(cur_dispatch_fiber_ != nullptr);
if (events.expire_at_ > 0) {
expire_events_[events.expire_at_].insert(cur_dispatch_fiber_);
cur_dispatch_fiber_->SetExpireEvent(events);
LOG_DEBUG("register fiber [%lu] with expire event at %ld", cur_dispatch_fiber_->Seq(), events.expire_at_);
}
for (size_t i = 0; i < events.waiting_fds_r_.size(); i++) {
int fd = events.waiting_fds_r_[i];
auto iter = io_waiting_fibers_.find(fd);
if (iter == io_waiting_fibers_.end()) {
io_waiting_fibers_.insert(std::make_pair(fd, WaitingFibers(cur_dispatch_fiber_, nullptr)));
cur_dispatch_fiber_->SetReadEvent(events);
}
}
for (size_t i = 0; i < events.waiting_fds_w_.size(); i++) {
int fd = events.waiting_fds_w_[i];
auto iter = io_waiting_fibers_.find(fd);
if (iter == io_waiting_fibers_.end()) {
io_waiting_fibers_.insert(std::make_pair(fd, WaitingFibers(nullptr, cur_dispatch_fiber_)));
cur_dispatch_fiber_->SetWriteEvent(events);
}
}
for (size_t i = 0; i < events.waiting_sems_.size(); i++) {
Sem &sem = events.waiting_sems_[i];
if (sem_infos_.find(sem) == sem_infos_.end()) {
sem_infos_.insert(std::make_pair(sem, SemInfo()));
}
sem_infos_[sem].fibers_.insert(cur_dispatch_fiber_);
}
}
//>提示<: 第一版 性能大概在3w~ 但是 会进行收尾工作. 只要是手动析构的,会进行最后一步收尾工作, 逻辑上更好!
// 将fd 取消注册
//>重点<: 知识取消了epoll注册,和各种的映射关系, 还需要最后一次上cpu 执行自己最后的处理逻辑
/*
bool XFiber::UnregisterFd(int fd) {
auto iter = io_waiting_fibers_.find(fd);
// 如果存在于 io_waitiing_fibers_ 映射中
if (iter != io_waiting_fibers_.end()) {
auto &waiting_fibers = iter->second;
Fiber *fiber_r = waiting_fibers.r_;
Fiber *fiber_w = waiting_fibers.w_;
// 如果读 fd映射的负责读fiber不为空
if (fiber_r != nullptr) {
WaitingEvents &evs_r = fiber_r->GetWaitingEvents();
for (size_t i = 0; i < evs_r.waiting_fds_r_.size(); i++) {
if (evs_r.waiting_fds_r_[i] == fd) {
int64_t expired_at = evs_r.expire_at_; // 获取超时时间
if (expired_at > 0) { //存在超时时间
auto expired_iter = expire_events_.find(expired_at);
// 断开超时队列 映射关系
if (expired_iter->second.find(fiber_r) == expired_iter->second.end()) {
LOG_ERROR("not fiber [%lu] in expired events", fiber_r->Seq());
}
else {
// 断开超时映射关系
expired_iter->second.erase(fiber_r);
}
}
}
}
}
// 读事件 超时队列断开映射
if (fiber_w != nullptr) {
WaitingEvents &evs_w = fiber_w->GetWaitingEvents();
for (size_t i = 0; i < evs_w.waiting_fds_w_.size(); i++) {
if (evs_w.waiting_fds_w_[i] == fd) {
int64_t expired_at = evs_w.expire_at_;
if (expired_at > 0) {
auto expired_iter = expire_events_.find(expired_at);
if (expired_iter->second.find(fiber_w) == expired_iter->second.end()) {
LOG_ERROR("not fiber [%lu] in expired events", fiber_w->Seq());
}
else {
expired_iter->second.erase(fiber_r);
}
}
}
}
}
io_waiting_fibers_.erase(iter);
}
else {
LOG_INFO("fd[%d] not register into sched", fd);
}
// 从epoll中取消注册
struct epoll_event ev;
ev.events = EPOLLIN | EPOLLOUT | EPOLLET;
ev.data.fd = fd;
if (epoll_ctl(ep_fd_, EPOLL_CTL_DEL, fd, &ev) < 0) {
LOG_ERROR("unregister fd[%d] from epoll efd[%d] failed, msg=%s", fd, ep_fd_, strerror(errno));
}
else {
LOG_INFO("unregister fd[%d] from epoll efd[%d] success!", fd, ep_fd_);
}
return true;
}
*/
//>第二版<: 会进行最后一次wakeup 收尾 更好 性能 3w+
bool XFiber::UnregisterFd(int fd) {
LOG_DEBUG("unregister fd[%d] from sheduler", fd);
auto io_waiting_fibers_iter = io_waiting_fibers_.find(fd);
// assert(io_waiting_fibers_iter != io_waiting_fibers_.end());
if (io_waiting_fibers_iter != io_waiting_fibers_.end()) {
WaitingFibers &waiting_fibers = io_waiting_fibers_iter->second;
if (waiting_fibers.r_ != nullptr) {
//为什么这里是唤醒?
//>重点<: wakeup 内部 会自动的断开其余的多余映射 and 处理 超时队列 (提前唤醒)
WakeUpFiber(waiting_fibers.r_);
}
if (waiting_fibers.w_ != nullptr) {
WakeUpFiber(waiting_fibers.w_);
}
io_waiting_fibers_.erase(io_waiting_fibers_iter);
}
struct epoll_event ev;
if (epoll_ctl(ep_fd_, EPOLL_CTL_DEL, fd, &ev) < 0) {
LOG_ERROR("unregister fd[%d] from epoll efd[%d] failed, msg=%s", fd, ep_fd_, strerror(errno));
}
else {
LOG_INFO("unregister fd[%d] from epoll efd[%d] success!", fd, ep_fd_);
}
return true;
}
//功能: 将事件集合 注册给对应fiber对象
// >重点<: 超时时间针对的是一整个fiber
//>提示<: 超时时间,是整个fiber协程对象的生存线
// 注册时, 需要将 fiber | 超时界限 注册到 全局超时队列里面去
// 同时 需要 根据传进来的 waitingevent 来分别初始化 fiber的 读写事件
// void XFiber::RegisterWaitingEvents(WaitingEvents &events) {
// // assert(curr_fiber_ != nullptr);
// assert(cur_dispatch_fiber_ != nullptr);
// //如果有超时时间设置: 将超时时间加入到超时队列并设置fiber的waitingevents
// // 超时队列: 超时时间,<fiber*,fiber*> 映射
// if (events.expire_at_ > 0) {
// // expire_events_[events.expire_at_].insert(curr_fiber_);
// //超时事件集合中 加入 当前 已经设置里超时时间的 fiber
// // 调度器的全局超时队列中加入当前正在运行的fiber
// expire_events_[events.expire_at_].insert(cur_dispatch_fiber_);
// // 初始化当前fiber的 waitingevents
// cur_dispatch_fiber_->SetWaitingEvent(events);
// // curr_fiber_->SetWaitingEvent(events);
// LOG_DEBUG("register fiber [%lu] with expire event at %ld", cur_dispatch_fiber_->Seq(), events.expire_at_);
// }
// //处理读事件: 将读事件 加入到等待集合 并设置fiber的waitingevents
// for (size_t i = 0; i < events.waiting_fds_r_.size(); i++) {
// int fd = events.waiting_fds_r_[i];
// auto iter = io_waiting_fibers_.find(fd);
// if (iter == io_waiting_fibers_.end()) {
// // 所有的 读事件 都加入到 等待集合epoll 监听映射的map 中去
// // pair: <fd,(r,w)> : 即 当前监听目标fd 和 当前fiber())
// io_waiting_fibers_.insert(std::make_pair(fd, WaitingFibers(cur_dispatch_fiber_, nullptr)));
// // 分别设置 当前fiber的 读事件
// cur_dispatch_fiber_->SetWaitingEvent(events);
// }
// }
// //处理写事件: 将写事件 加入到等待集合 并设置fiber的waitingevents
// for (size_t i = 0; i < events.waiting_fds_w_.size(); i++) {
// int fd = events.waiting_fds_w_[i];
// auto iter = io_waiting_fibers_.find(fd);
// if (iter == io_waiting_fibers_.end()) {
// // 所有的 写事件都加入到 等待集合map中去
// io_waiting_fibers_.insert(std::make_pair(fd, WaitingFibers(nullptr, cur_dispatch_fiber_)));
// //分别设置当前fiber自身的 waitingevent 中的写事件
// cur_dispatch_fiber_->SetWaitingEvent(events);
// }
// }
// //处理信号量: 将信号量加入到信号集合 waitingevents
// for (size_t i = 0; i < events.waiting_sems_.size(); i++) {
// Sem &sem = events.waiting_sems_[i];
// // 没有这个信号量存在,则新插入
// if (sem_infos_.find(sem) == sem_infos_.end()) {
// sem_infos_.insert(std::make_pair(sem, SemInfo()));
// }
// //>重点<: 为什么这里和生面的不一样?
// //>解答<:
// sem_infos_[sem].fibers_.insert(cur_dispatch_fiber_);
// }
// }
// map: fd waitingfibers
// waitingfibers : fiber1* fiber2*
// waitingevents: 读数组 写数组 信号量数组 超时时间
// >重点<: 根据 传入的 waitingevents 来分别设置 fiber的 读写事件 (读事件 写事件分开))
// >提示<: 传入的events 是临时变量 ,fiber内部也有一个和这个一样结构的events事件集合
// 功能: 将events内部的事件 注册到fd中 (超时事件,r_fd,w_fd))
// void Fiber::SetWaitingEvent(const WaitingEvents &events) {
// //注册读事件
// for (size_t i = 0; i < events.waiting_fds_r_.size(); i++) {
// //fiber 的waitingevents
// this->waiting_events_.waiting_fds_r_.push_back(events.waiting_fds_r_[i]);
// }
// // 注册写事件
// for (size_t i = 0; i < events.waiting_fds_w_.size(); i++) {
// this->waiting_events_.waiting_fds_w_.push_back(events.waiting_fds_w_[i]);
// }
// // 设置超时时间: 就是整个fiber对象的 生存时间线
// if (events.expire_at_ > 0) {
// this->waiting_events_.expire_at_ = events.expire_at_;
// }
// }
// 在epoll中注册某个fd
void XFiber::TakeOver(int fd) {
struct epoll_event ev;
// 读写 ET
ev.events = EPOLLIN | EPOLLOUT | EPOLLET;
ev.data.fd = fd;// 联合体 union 之关心一个
/* epoll_ctl 函数说明: lepoll的事件注册函数
int epoll_ctl(int epfd, int op, int fd, struct epoll_event *event)
参数:
epfd: epoll结构句柄
op: epoll_ctl注册的动作
1. EPOLL_CTL_ADD:注册新的fd到epfd中;
2. EPOLL_CTL_MOD:修改已经注册的fd的监听事件;
3. EPOLL_CTL_DEL:从epfd中删除一个fd;
fd : 需要epoll结构 去侦听的 句柄
event: 关心 被侦听fd 上发生的事件类型 (只关心这些类型,其他的不关心)
EPOLLIN : 读事件
EPOLLOUT: 写事件
EPOLLPRI: 额外数据到来(紧急数据)
EPOLLERR:发生错误
EPOLLET: 表示使用ET模式
返回值:
成功: 0
失败: -1
*/
if (epoll_ctl(ep_fd_, EPOLL_CTL_ADD, fd, &ev) < 0) {
LOG_ERROR("add fd [%d] into epoll failed, msg=%s", fd, strerror(errno));
exit(-1);
}
LOG_DEBUG("add fd[%d] into epoll event success", fd);
}
/*
//从epoll中取消注册某个fd
//>提问<: 为什么不在超时队列里面也删除?
//>解答<: 从epoll中删除 就不会被触发了 等到超时,他自己就会wakeup 之后自己自觉注销
bool XFiber::UnregisterFd(int fd) {
LOG_DEBUG("unregister fd[%d] from sheduler", fd);
auto io_waiting_fibers_iter = io_waiting_fibers_.find(fd);
// assert(io_waiting_fibers_iter != io_waiting_fibers_.end());
//第一步: 从等待集合map中删除 该fd对应的 读写fibers ==> fd都不见听了 fiber也就没用了
if (io_waiting_fibers_iter != io_waiting_fibers_.end()) {
WaitingFibers &waiting_fibers = io_waiting_fibers_iter->second;
// 删除之前,需要先 将他所负责的事情 进行交接 该干的要干完
if (waiting_fibers.r_ != nullptr) {
WakeUpFiber(waiting_fibers.r_);
}
if (waiting_fibers.w_ != nullptr) {
WakeUpFiber(waiting_fibers.w_);
}
io_waiting_fibers_.erase(io_waiting_fibers_iter);
}
struct epoll_event ev;
//第二步: 从epoll中删除
if (epoll_ctl(ep_fd_, EPOLL_CTL_DEL, fd, &ev) < 0) {
LOG_ERROR("unregister fd[%d] from epoll efd[%d] failed, msg=%s", fd, ep_fd_, strerror(errno));
}
else {
LOG_INFO("unregister fd[%d] from epoll efd[%d] success!", fd, ep_fd_);
}
//>提问<: 为什么 不去设置超时队列和信号量队列呢? 超时队列不用管,超时的fiber不会处理,并且超时队列最后会彻底clear
//>提问<: 那么信号量队列呢? epoll主要关心的还是读写
return true;
}
*/
//>注意<: 这一居有什么用? 为什么这么写? 和SEQ有什么关系?
//>解答<: 可能和 网络通信阶段的 三握的seq油管
thread_local uint64_t fiber_seq = 0;
/* 函数参数解析
1. fiber对象(暨背调度的协程)的入口函数地址
2. 该入口函数地址是保存在 uc_mcontext中的:
3. uc_mcontext:保存所有上下文信息,寄存器信息,入口函数地址
4. stack_size:制定站空间多钱啊小
5. fiber_name: 当前协程 name
6. xfiber: 调度器
*/
Fiber::Fiber(std::function<void ()> run, XFiber *xfiber, size_t stack_size, std::string fiber_name) {
this->entry_func_ = run;
this->xfiber_ = xfiber;
this->fiber_name_ = fiber_name;
//有栈协程: 栈大小为 128kb
this->stack_size_ = stack_size;
//协程的栈起始位置
this->stack_ptr_ = new uint8_t[stack_size_];
//ucontext函数簇 ctx初始化
getcontext(&this->ctx_);
//==========三兄弟=========
ctx_.uc_stack.ss_sp = stack_ptr_;
ctx_.uc_stack.ss_size = stack_size_;
//>注意<: 这里很重要
//>重点<: 每一个 fiber的 ctx 都link上 xfiber的唯一的 ctx
//>提问<: 为什么要这样做?
//>解答<: 这样能够保证 任何一个 fiber 运行完毕后,都会自动的去切换到xfiber 然后继续有xfiber调度其他fiber
//---> 这也是协程调度器能够持续运行的根本
ctx_.uc_link = xfiber->SchedCtx();
//=========================
//该协程的入口函数 (makecontext会绑定协程的入口函数地址并且,进行上下文切换)
/*
makecontext参数解析:
1. &this->ctx_: 当前的上下文信息
2. Start: 协程入口函数地址
3. "1": 后续参数个数
4. "this": this指针
*/
makecontext(&this->ctx_, (void (*)())Fiber::Start, 1, this);
seq_ = fiber_seq;
fiber_seq++;
//makecontext成功之后, xfiber 单线程主协程被替换成为一个崭新的协程
// so 当前fiber的状态时 INIT
status_ = FiberStatus::INIT;
}
//>注意<: 销毁一个协程
//需要销毁一个协程的栈空间
Fiber::~Fiber() {
// 删除 fiber的携程栈
delete []stack_ptr_;
stack_ptr_ = nullptr;
stack_size_ = 0;
}
uint64_t Fiber::Seq() {
return this->seq_;
}
// 获取 当前fiber 的上下文信息
XFiberCtx *Fiber::Get_Fiber_Ctx() {
return &ctx_;
}
void Fiber::Start(Fiber * fiber) {
// 在这里 执行 fiber真正的处理函数也就是入口函数
fiber->entry_func_(); // 在这里之后 陷入执行逻辑
// 处理完成 fiber彻底结束 ,切换状态
fiber->status_ = FiberStatus::FINISHED;
LOG_DEBUG("fiber[%lu] finished...", fiber->Seq());
}
std::string Fiber::Name() {
return fiber_name_;
}
bool Fiber::IsFinished() {
return status_ == FiberStatus::FINISHED;
}
// 这里是信号量
// struct Sem {
//>注意<: thread_local: 确保每一个县城里面 只有一个sem_seq
// thread_local static int64_t sem_seq;
// int64_t seq_;
// Sem(int32_t value);
// bool Acquire(int32_t apply_value, int32_t timeout_ms=-1);
// void Release(int32_t acquired_value);
// bool operator < (const Sem &other) const {
// return seq_ < other.seq_;
// }
// };
//>注意<: 信号量结构体
// struct SemInfo {
// SemInfo(int32_t value = 0) {
// value_ = value;
// UCBA_fiber_ = nullptr;
// UCBA_value_ = 0;
// fibers_.clear();
// }
// int32_t value_;