/
processor.go
149 lines (148 loc) · 3.45 KB
/
processor.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
package mongodb
//type MongoProcessor struct {
// config MongoProcessorConfig
//
// readPool *Pool
// writePool *Pool
//}
//type MongoProcessorConfig struct {
// IdleConnectionTimeout time.Duration
//}
//
//func (m *MongoProcessor) Stop() {
// logrus.Info("MongoProcessor stopped")
//}
//
//func (m *MongoProcessor) Start() {
// logrus.Info("MongoProcessor started")
// // start consuming queue
//}
//
//func NewMongoProcessor(config MongoProcessorConfig) *MongoProcessor {
// //TODO move mongo url into config
// url := "172.28.152.101:27017"
//
// return &MongoProcessor{
// config: config,
// readPool: NewPool(url, 10),
// writePool: NewPool(url, 10),
// }
//}
//
//func (m *MongoProcessor) ProcessConnection(conn net.Conn) error {
// defer conn.Close()
//
// fmt.Println("start process connection")
//
// // http://docs.mongodb.org/manual/faq/diagnostics/#faq-keepalive
// if conn, ok := conn.(*net.TCPConn); ok {
// conn.SetKeepAlivePeriod(2 * time.Minute)
// conn.SetKeepAlive(true)
// }
//
// reader := bufio.NewReader(conn)
//
// backend := m.writePool.Acquire()
// defer m.writePool.Release(backend)
//
// for {
// conn.SetReadDeadline(time.Now().Add(m.config.IdleConnectionTimeout))
//
// cmdHeader := make([]byte, message.HeaderLen)
// _, err := reader.Read(cmdHeader)
// if err != nil {
// if err == io.EOF {
// fmt.Println("target closed")
// logrus.Info("target closed")
// return nil
// } else if neterr, ok := err.(net.Error); ok && neterr.Timeout() {
// fmt.Println("target timeout")
// logrus.Info("target timeout")
// conn.Close()
// return nil
// }
// return err
// }
//
// // query command
// msgSize := bytes.GetInt32(cmdHeader, 0)
// fmt.Println("msgsize: ", msgSize)
//
// cmdBody := make([]byte, msgSize-message.HeaderLen)
// _, err = reader.Read(cmdBody)
// if err != nil {
// fmt.Println("read body error: ", err)
// return err
// }
// fmt.Println(fmt.Sprintf("msg header: %x", cmdHeader))
// fmt.Println(fmt.Sprintf("msg body: %x", cmdBody))
//
// cmdFull := append(cmdHeader, cmdBody...)
// err = m.messageHandler(cmdFull, conn, backend)
// if err != nil {
// // TODO handle err
// return err
// }
//
// }
// return nil
//}
//
//func (m *MongoProcessor) messageHandler(bytes []byte, client, backend net.Conn) error {
// //
// //var msg message.RequestMessage
// //fmt.Println("Request--->")
// //fmt.Println(hex.Dump(bytes))
// //err := msg.Decode(bytes)
// //if err != nil {
// // // TODO handle err
// // return err
// //}
// //
// ////var pool *Pool
// ////if msg.ReadOnly() {
// //// pool = m.readPool
// ////} else {
// //// pool = m.writePool
// ////}
// ////backend := pool.Acquire()
// ////defer pool.Release(backend)
// //
// //err = msg.WriteTo(backend)
// //if err != nil {
// // // TODO handle err
// // return err
// //}
// //
// //var msgResp message.ResponseMessage
// //err = msgResp.ReadFromMongo(backend)
// //if err != nil {
// // // TODO handle err
// // return err
// //}
// //fmt.Println("<---Response")
// ////fmt.Println(hex.Dump(msgResp.payload))
// //err = msgResp.WriteTo(client)
// //if err != nil {
// // // TODO handle err
// // return err
// //}
// //
// ////err = m.handleBlockDBEvents(&msgResp)
// ////if err != nil {
// //// // TODO handle err
// //// return err
// ////}
//
// return nil
//}
//
//func (m *MongoProcessor) handleBlockDBEvents(msg message.Message) error {
// // TODO not implemented yet
//
// events := msg.ParseCommand()
//
// fmt.Println("block db events: ", events)
//
// return nil
//}