forked from Allenxuxu/gev
-
Notifications
You must be signed in to change notification settings - Fork 0
/
server.go
139 lines (116 loc) · 3.03 KB
/
server.go
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
package gev
import (
"errors"
"runtime"
"time"
"github.com/Allenxuxu/gev/connection"
"github.com/Allenxuxu/gev/eventloop"
"github.com/Allenxuxu/gev/listener"
"github.com/Allenxuxu/gev/log"
"github.com/Allenxuxu/toolkit/sync"
"github.com/Allenxuxu/toolkit/sync/atomic"
"github.com/RussellLuo/timingwheel"
"golang.org/x/sys/unix"
)
// Handler Server 注册接口
type Handler interface {
connection.CallBack
OnConnect(c *connection.Connection)
}
// Server gev Server
type Server struct {
loop *eventloop.EventLoop
workLoops []*eventloop.EventLoop
callback Handler
timingWheel *timingwheel.TimingWheel
opts *Options
running atomic.Bool
}
// NewServer 创建 Server
func NewServer(handler Handler, opts ...Option) (server *Server, err error) {
if handler == nil {
return nil, errors.New("handler is nil")
}
options := newOptions(opts...)
server = new(Server)
server.callback = handler
server.opts = options
server.timingWheel = timingwheel.NewTimingWheel(server.opts.tick, server.opts.wheelSize)
server.loop, err = eventloop.New()
if err != nil {
_ = server.loop.Stop()
return nil, err
}
l, err := listener.New(server.opts.Network, server.opts.Address, options.ReusePort, server.loop, server.handleNewConnection)
if err != nil {
return nil, err
}
if err = server.loop.AddSocketAndEnableRead(l.Fd(), l); err != nil {
return nil, err
}
if server.opts.NumLoops <= 0 {
server.opts.NumLoops = runtime.NumCPU()
}
wloops := make([]*eventloop.EventLoop, server.opts.NumLoops)
for i := 0; i < server.opts.NumLoops; i++ {
l, err := eventloop.New()
if err != nil {
for j := 0; j < i; j++ {
_ = wloops[j].Stop()
}
return nil, err
}
wloops[i] = l
}
server.workLoops = wloops
return
}
// RunAfter 延时任务
func (s *Server) RunAfter(d time.Duration, f func()) *timingwheel.Timer {
return s.timingWheel.AfterFunc(d, f)
}
// RunEvery 定时任务
func (s *Server) RunEvery(d time.Duration, f func()) *timingwheel.Timer {
return s.timingWheel.ScheduleFunc(&everyScheduler{Interval: d}, f)
}
func (s *Server) handleNewConnection(fd int, sa unix.Sockaddr) {
loop := s.opts.Strategy(s.workLoops)
c := connection.New(fd, loop, sa, s.opts.Protocol, s.timingWheel, s.opts.IdleTime, s.callback)
loop.QueueInLoop(func() {
s.callback.OnConnect(c)
if err := loop.AddSocketAndEnableRead(fd, c); err != nil {
log.Error("[AddSocketAndEnableRead]", err)
}
})
}
// Start 启动 Server
func (s *Server) Start() {
sw := sync.WaitGroupWrapper{}
s.timingWheel.Start()
length := len(s.workLoops)
for i := 0; i < length; i++ {
sw.AddAndRun(s.workLoops[i].RunLoop)
}
sw.AddAndRun(s.loop.RunLoop)
s.running.Set(true)
sw.Wait()
}
// Stop 关闭 Server
func (s *Server) Stop() {
if s.running.Get() {
s.running.Set(false)
s.timingWheel.Stop()
if err := s.loop.Stop(); err != nil {
log.Error(err)
}
for k := range s.workLoops {
if err := s.workLoops[k].Stop(); err != nil {
log.Error(err)
}
}
}
}
// Options 返回 options
func (s *Server) Options() Options {
return *s.opts
}