-
Notifications
You must be signed in to change notification settings - Fork 2
/
zconnection.go
179 lines (151 loc) · 3.76 KB
/
zconnection.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
package server
import (
"net"
"io"
"log"
"errors"
"github.com/ZhengjunHUO/zjunx/pkg/encoding"
)
type ZConnection interface {
Start()
Reader()
Writer()
Close()
RespondToClient(encoding.ZContentType, []byte) error
GetID() uint64
GetServer() ZServer
UpdateContext(string, interface{})
GetContext(string) interface{}
DeleteContext(string)
CallPostStart()
CallPreStop()
}
type Connection struct {
ID uint64
Conn net.Conn
Server ZServer
Context map[string]interface{}
chServerResp chan []byte
chClose chan bool
isActive bool
}
func ConnInit(cnxID uint64, conn net.Conn, s ZServer) ZConnection {
cnx := &Connection{
ID: cnxID,
Conn: conn,
Server: s,
Context: make(map[string]interface{}),
chServerResp: make(chan []byte),
chClose: make(chan bool, 1),
isActive: true,
}
s.GetCnxAdm().Register(cnx)
return cnx
}
// Read from the TCP stream payload and decode the raw bytes to struct
// Prepare a processed request and send it to a worker to handle it
func (c *Connection) Reader() {
defer c.Close()
defer c.Server.GetCnxAdm().Remove(c)
blk := encoding.BlockInit()
for {
ct := encoding.ContentInit(encoding.ZContentType(0), []byte{})
if err := blk.Unmarshalling(c.Conn, ct); err != nil {
if err != io.EOF {
log.Println("[WARN] Unmarshalling failed: ", err)
}
break
}
// build a request from a valid incoming package
req := ReqInit(c, ct)
// sending request to a worker
c.Server.GetMux().Schedule(req)
log.Printf("[DEBUG] Request from Connection [id: %v] sheduled.\n", c.GetID())
}
}
// Called by handler after dealing with the request
// Send the raw bytes to Writer
func (c *Connection) RespondToClient(ct encoding.ZContentType, data []byte) error {
if ! c.isActive {
return errors.New("Error sending response: Connection is closed")
}
blk := encoding.BlockInit()
buf, err := blk.Marshalling(encoding.ContentInit(ct, data))
if err != nil {
return errors.New("Error sending response to client !")
}
// Writer process listening at the other end of channel
c.chServerResp <- buf
return nil
}
// Write the processed response (raw bytes received from handler) to client
// Quit on receiving the close signal after Reader quits
func (c *Connection) Writer() {
for {
select {
case data := <- c.chServerResp:
if _, err := c.Conn.Write(data); err != nil {
log.Println("[ERROR] Write back to client: ", err)
return
}
case <- c.chClose:
return
}
}
}
// Seperate read/write thread, leave the handling part to ZMux
func (c *Connection) Start() {
log.Printf("[DEBUG] Connection [id: %d] established from %v\n", c.ID, c.Conn.RemoteAddr())
go c.Reader()
go c.Writer()
c.CallPostStart()
}
func (c *Connection) GetID() uint64 {
return c.ID
}
func (c *Connection) GetServer() ZServer {
return c.Server
}
func (c *Connection) UpdateContext(key string, value interface{}) {
c.Context[key] = value
}
func (c *Connection) GetContext(key string) interface{} {
if v, ok := c.Context[key]; ok {
return v
}else{
return nil
}
}
func (c *Connection) DeleteContext(key string) {
delete(c.Context, key)
}
func (c *Connection) CallPostStart() {
if c.Server.GetPostStartHook() == nil {
return
}
c.Server.GetPostStartHook()(c)
}
func (c *Connection) CallPreStop() {
if c.Server.GetPreStopHook() == nil {
return
}
c.Server.GetPreStopHook()(c)
}
// Clenup current connection before exit
func (c *Connection) Close() {
log.Printf("[DEBUG] Closing connection [id: %d] ... \n", c.ID)
c.CallPreStop()
if !c.isActive {
return
}
c.isActive = false
// close the Writer process
c.chClose <- true
close(c.chClose)
close(c.chServerResp)
if err := c.Conn.Close(); err !=nil {
log.Printf("[DEBUG] Closing connection [id: %d]: %s\n", c.ID, err)
}else{
log.Printf("[DEBUG] Connection [id: %d] closed. \n", c.ID)
}
}