-
Notifications
You must be signed in to change notification settings - Fork 62
/
node_tcpros_server.go
85 lines (72 loc) · 1.49 KB
/
node_tcpros_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
package goroslib
import (
"net"
"sync"
"time"
"github.com/bluenviron/goroslib/v2/pkg/protocommon"
"github.com/bluenviron/goroslib/v2/pkg/prototcp"
)
func (n *Node) runTcprosServer(wg *sync.WaitGroup) {
defer wg.Done()
for {
conn, err := n.tcprosListener.Accept()
if err != nil {
break
}
select {
case n.tcpConnNew <- conn:
case <-n.ctx.Done():
conn.Close()
}
}
}
func (n *Node) runTcprosServerConn(wg *sync.WaitGroup, nconn net.Conn) {
defer wg.Done()
ok := func() bool {
tconn := prototcp.NewConn(nconn)
nconn.SetReadDeadline(time.Now().Add(n.conf.ReadTimeout))
rawHeader, err := tconn.ReadHeaderRaw()
if err != nil {
return false
}
if _, ok := rawHeader["topic"]; ok {
var header prototcp.HeaderSubscriber
err = protocommon.HeaderDecode(rawHeader, &header)
if err != nil {
return false
}
select {
case n.tcpConnSubscriber <- tcpConnSubscriberReq{
nconn: nconn,
tconn: tconn,
header: &header,
}:
case <-n.ctx.Done():
}
return true
} else if _, ok := rawHeader["service"]; ok {
var header prototcp.HeaderServiceClient
err = protocommon.HeaderDecode(rawHeader, &header)
if err != nil {
return false
}
select {
case n.tcpConnServiceClient <- tcpConnServiceClientReq{
nconn: nconn,
tconn: tconn,
header: &header,
}:
case <-n.ctx.Done():
}
return true
}
return false
}()
if !ok {
nconn.Close()
select {
case n.tcpConnClose <- nconn:
case <-n.ctx.Done():
}
}
}