/
websocketController.go
351 lines (298 loc) · 9.05 KB
/
websocketController.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
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
package web
import (
"sync"
"github.com/henrylee2cn/pholcus/app"
"github.com/henrylee2cn/pholcus/app/spider"
"github.com/henrylee2cn/pholcus/common/util"
ws "github.com/henrylee2cn/pholcus/common/websocket"
"github.com/henrylee2cn/pholcus/config"
"github.com/henrylee2cn/pholcus/logs"
"github.com/henrylee2cn/pholcus/runtime/status"
)
type SocketController struct {
connPool map[string]*ws.Conn
wchanPool map[string]*Wchan
connRWMutex sync.RWMutex
wchanRWMutex sync.RWMutex
}
func (self *SocketController) GetConn(sessID string) *ws.Conn {
self.connRWMutex.RLock()
defer self.connRWMutex.RUnlock()
return self.connPool[sessID]
}
func (self *SocketController) GetWchan(sessID string) *Wchan {
self.wchanRWMutex.RLock()
defer self.wchanRWMutex.RUnlock()
return self.wchanPool[sessID]
}
func (self *SocketController) Add(sessID string, conn *ws.Conn) {
self.connRWMutex.Lock()
self.wchanRWMutex.Lock()
defer self.connRWMutex.Unlock()
defer self.wchanRWMutex.Unlock()
self.connPool[sessID] = conn
self.wchanPool[sessID] = newWchan()
}
func (self *SocketController) Remove(sessID string, conn *ws.Conn) {
self.connRWMutex.Lock()
self.wchanRWMutex.Lock()
defer self.connRWMutex.Unlock()
defer self.wchanRWMutex.Unlock()
if self.connPool[sessID] == nil {
return
}
wc := self.wchanPool[sessID]
close(wc.wchan)
conn.Close()
delete(self.connPool, sessID)
delete(self.wchanPool, sessID)
}
func (self *SocketController) Write(sessID string, void map[string]interface{}, to ...int) {
self.wchanRWMutex.RLock()
defer self.wchanRWMutex.RUnlock()
// to为1时,只向当前连接发送;to为-1时,向除当前连接外的其他所有连接发送;to为0时或为空时,向所有连接发送
var t int = 0
if len(to) > 0 {
t = to[0]
}
void["mode"] = app.LogicApp.GetAppConf("mode").(int)
switch t {
case 1:
wc := self.wchanPool[sessID]
if wc == nil {
return
}
void["initiative"] = true
wc.wchan <- void
case 0, -1:
l := len(self.wchanPool)
for _sessID, wc := range self.wchanPool {
if t == -1 && _sessID == sessID {
continue
}
_void := make(map[string]interface{}, l)
for k, v := range void {
_void[k] = v
}
if _sessID == sessID {
_void["initiative"] = true
} else {
_void["initiative"] = false
}
wc.wchan <- _void
}
}
}
type Wchan struct {
wchan chan interface{}
}
func newWchan() *Wchan {
return &Wchan{
wchan: make(chan interface{}, 1024),
}
}
var (
wsApi = map[string]func(string, map[string]interface{}){}
Sc = &SocketController{
connPool: make(map[string]*ws.Conn),
wchanPool: make(map[string]*Wchan),
}
)
func wsHandle(conn *ws.Conn) {
defer func() {
if p := recover(); p != nil {
logs.Log.Error("%v", p)
}
}()
sess, _ := globalSessions.SessionStart(nil, conn.Request())
sessID := sess.SessionID()
if Sc.GetConn(sessID) == nil {
Sc.Add(sessID, conn)
}
defer Sc.Remove(sessID, conn)
go func() {
var err error
for info := range Sc.GetWchan(sessID).wchan {
if _, err = ws.JSON.Send(conn, info); err != nil {
return
}
}
}()
for {
var req map[string]interface{}
if err := ws.JSON.Receive(conn, &req); err != nil {
// logs.Log.Debug("websocket接收出错断开 (%v) !", err)
return
}
// log.Log.Debug("Received from web: %v", req)
wsApi[util.Atoa(req["operate"])](sessID, req)
}
}
func init() {
// 初始化运行
wsApi["refresh"] = func(sessID string, req map[string]interface{}) {
// 写入发送通道
Sc.Write(sessID, tplData(app.LogicApp.GetAppConf("mode").(int)), 1)
}
// 初始化运行
wsApi["init"] = func(sessID string, req map[string]interface{}) {
var mode = util.Atoi(req["mode"])
var port = util.Atoi(req["port"])
var master = util.Atoa(req["ip"]) //服务器(主节点)地址,不含端口
currMode := app.LogicApp.GetAppConf("mode").(int)
if currMode == status.UNSET {
app.LogicApp.Init(mode, port, master, Lsc) // 运行模式初始化,设置log输出目标
} else {
app.LogicApp = app.LogicApp.ReInit(mode, port, master) // 切换运行模式
}
if mode == status.CLIENT {
go app.LogicApp.Run()
}
// 写入发送通道
Sc.Write(sessID, tplData(mode))
}
wsApi["run"] = func(sessID string, req map[string]interface{}) {
if app.LogicApp.GetAppConf("mode").(int) != status.CLIENT {
setConf(req)
}
if app.LogicApp.GetAppConf("mode").(int) == status.OFFLINE {
Sc.Write(sessID, map[string]interface{}{"operate": "run"})
}
go func() {
app.LogicApp.Run()
if app.LogicApp.GetAppConf("mode").(int) == status.OFFLINE {
Sc.Write(sessID, map[string]interface{}{"operate": "stop"})
}
}()
}
// 终止当前任务,现仅支持单机模式
wsApi["stop"] = func(sessID string, req map[string]interface{}) {
if app.LogicApp.GetAppConf("mode").(int) != status.OFFLINE {
Sc.Write(sessID, map[string]interface{}{"operate": "stop"})
return
} else {
// println("stopping^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^")
app.LogicApp.Stop()
// println("stopping++++++++++++++++++++++++++++++++++++++++")
Sc.Write(sessID, map[string]interface{}{"operate": "stop"})
// println("stopping$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$$")
}
}
// 任务暂停与恢复,目前仅支持单机模式
wsApi["pauseRecover"] = func(sessID string, req map[string]interface{}) {
if app.LogicApp.GetAppConf("mode").(int) != status.OFFLINE {
return
}
app.LogicApp.PauseRecover()
Sc.Write(sessID, map[string]interface{}{"operate": "pauseRecover"})
}
// 退出当前模式
wsApi["exit"] = func(sessID string, req map[string]interface{}) {
app.LogicApp = app.LogicApp.ReInit(status.UNSET, 0, "")
Sc.Write(sessID, map[string]interface{}{"operate": "exit"})
}
}
func tplData(mode int) map[string]interface{} {
var info = map[string]interface{}{"operate": "init", "mode": mode}
// 运行模式标题
switch mode {
case status.OFFLINE:
info["title"] = config.FULL_NAME + " 【 运行模式 -> 单机 】"
case status.SERVER:
info["title"] = config.FULL_NAME + " 【 运行模式 -> 服务端 】"
case status.CLIENT:
info["title"] = config.FULL_NAME + " 【 运行模式 -> 客户端 】"
}
if mode == status.CLIENT {
return info
}
// 蜘蛛家族清单
info["spiders"] = map[string]interface{}{
"menu": spiderMenu,
"curr": func() interface{} {
l := app.LogicApp.GetSpiderQueue().Len()
if l == 0 {
return 0
}
var curr = make(map[string]bool, l)
for _, sp := range app.LogicApp.GetSpiderQueue().GetAll() {
curr[sp.GetName()] = true
}
return curr
}(),
}
// 输出方式清单
info["OutType"] = map[string]interface{}{
"menu": app.LogicApp.GetOutputLib(),
"curr": app.LogicApp.GetAppConf("OutType"),
}
// 并发协程上限
info["ThreadNum"] = map[string]int{
"max": 999999,
"min": 1,
"curr": app.LogicApp.GetAppConf("ThreadNum").(int),
}
// 暂停区间/ms(随机: Pausetime/2 ~ Pausetime*2)
info["Pausetime"] = map[string][]int64{
"menu": {0, 100, 300, 500, 1000, 3000, 5000, 10000, 15000, 20000, 30000, 60000},
"curr": []int64{app.LogicApp.GetAppConf("Pausetime").(int64)},
}
// 代理IP更换的间隔分钟数
info["ProxyMinute"] = map[string][]int64{
"menu": {0, 1, 3, 5, 10, 15, 20, 30, 45, 60, 120, 180},
"curr": []int64{app.LogicApp.GetAppConf("ProxyMinute").(int64)},
}
// 分批输出的容量
info["DockerCap"] = map[string]int{
"min": 1,
"max": 5000000,
"curr": app.LogicApp.GetAppConf("DockerCap").(int),
}
// 采集上限
if app.LogicApp.GetAppConf("Limit").(int64) == spider.LIMIT {
info["Limit"] = 0
} else {
info["Limit"] = app.LogicApp.GetAppConf("Limit")
}
// 自定义配置
info["Keyins"] = app.LogicApp.GetAppConf("Keyins")
// 继承历史记录
info["SuccessInherit"] = app.LogicApp.GetAppConf("SuccessInherit")
info["FailureInherit"] = app.LogicApp.GetAppConf("FailureInherit")
// 运行状态
info["status"] = app.LogicApp.Status()
return info
}
// 配置运行参数
func setConf(req map[string]interface{}) {
if tn := util.Atoi(req["ThreadNum"]); tn == 0 {
app.LogicApp.SetAppConf("ThreadNum", 1)
} else {
app.LogicApp.SetAppConf("ThreadNum", tn)
}
app.LogicApp.
SetAppConf("Pausetime", int64(util.Atoi(req["Pausetime"]))).
SetAppConf("ProxyMinute", int64(util.Atoi(req["ProxyMinute"]))).
SetAppConf("OutType", util.Atoa(req["OutType"])).
SetAppConf("DockerCap", util.Atoi(req["DockerCap"])).
SetAppConf("Limit", int64(util.Atoi(req["Limit"]))).
SetAppConf("Keyins", util.Atoa(req["Keyins"])).
SetAppConf("SuccessInherit", req["SuccessInherit"] == "true").
SetAppConf("FailureInherit", req["FailureInherit"] == "true")
setSpiderQueue(req)
}
func setSpiderQueue(req map[string]interface{}) {
spNames, ok := req["spiders"].([]interface{})
if !ok {
return
}
spiders := []*spider.Spider{}
for _, sp := range app.LogicApp.GetSpiderLib() {
for _, spName := range spNames {
if util.Atoa(spName) == sp.GetName() {
spiders = append(spiders, sp.Copy())
}
}
}
app.LogicApp.SpiderPrepare(spiders)
}