-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathnetstring.go
More file actions
102 lines (91 loc) · 2.2 KB
/
Copy pathnetstring.go
File metadata and controls
102 lines (91 loc) · 2.2 KB
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
package netString
import (
"bufio"
"fmt"
"net"
"strconv"
)
const maxMsg uint64 = 1024 * 8
type NetStringConn struct {
Conn net.Conn
MaxMsg uint64
queueOut chan string
}
type NetStringProcessor interface {
Connected(*NetStringConn)
Msg(*NetStringConn, string)
Disconected(*NetStringConn, error)
}
type NetString struct {
MaxMsg uint64
nsp NetStringProcessor
}
func newNetString(nsp NetStringProcessor) *NetString {
ns := NetString{maxMsg, nsp}
return &ns
}
func (n *NetString) Connect(conn net.Conn) {
n.process(conn)
}
func (n *NetString) Listen(ln net.Listener) error {
for {
conn, err := ln.Accept()
if err != nil {
return err
}
go n.process(conn)
}
}
func (n *NetString) process(conn net.Conn) {
nsc := &NetStringConn{conn, n.MaxMsg, make(chan string, 10)}
go nsc.process()
n.nsp.Connected(nsc)
reader := bufio.NewReader(conn)
for {
longs, err := reader.ReadString(':')
if err != nil {
n.nsp.Disconected(nsc, err)
break
}
long, err := strconv.ParseUint(longs[0:len(longs)-1], 10, 64)
if err != nil {
n.nsp.Disconected(nsc, fmt.Errorf("error in convert len data: %s\n", err.Error()))
break
}
fmt.Printf("max: %v\n", nsc.MaxMsg)
if long > nsc.MaxMsg {
n.nsp.Disconected(nsc, fmt.Errorf("Data len (%d) is bigger than max msg len (%d)\n", long, nsc.MaxMsg))
break
}
b := make([]byte, long)
if nread, err := reader.Read(b); err != nil {
n.nsp.Disconected(nsc, fmt.Errorf("error in receive data: %s\n", err.Error()))
break
} else if nread != len(b) {
n.nsp.Disconected(nsc, fmt.Errorf("error, receive %d data and expected %d\n", nread, long))
break
}
if by, err := reader.ReadByte(); err != nil {
n.nsp.Disconected(nsc, fmt.Errorf("error in last byte, must be ',': %s", err.Error()))
break
} else if by != ',' {
n.nsp.Disconected(nsc, fmt.Errorf("error, last byte must be ',' and is '%v'", by))
break
}
n.nsp.Msg(nsc, string(b))
}
close(nsc.queueOut)
}
func (n *NetStringConn) process() {
w := bufio.NewWriter(n.Conn)
for s := range n.queueOut {
w.Write([]byte(strconv.Itoa(len(s))))
w.WriteByte(':')
w.WriteString(s)
w.WriteByte(',')
w.Flush()
}
}
func (n *NetStringConn) Send(s string) {
n.queueOut <- s
}