@cryptotaxi247 / kubo / commits / c90c16eee

updated multiaddr + added utp

Juan Batiz-Benet committed Nov 20, 2014 at 05:10 UTC c90c16eee4bad254ffbc576e7d495b644c9350d7
23 files changed +3046 -11
Godeps/Godeps.json
+5 -1
@@ -81,6 +81,10 @@
81 "ImportPath": "github.com/gorilla/mux",
82 "Rev": "4b8fbc56f3b2400a7c7ea3dba9b3539787c486b6"
83 },
84 + {
85 + "ImportPath": "github.com/h2so5/utp",
86 + "Rev": "654d875bb65e96729678180215cf080fe2810371"
87 + },
88 {
89 "ImportPath": "github.com/inconshreveable/go-update",
90 "Rev": "221d034a558b4c21b0624b2a450c076913854a57"
@@ -112,7 +116,7 @@
116 },
117 {
118 "ImportPath": "github.com/jbenet/go-multiaddr-net",
115 - "Rev": "625fac6e5073702ac9cd6028b9dc2457fd9cbf9d"
119 + "Rev": "b6265d8119558acf3912db44abb34d97c30c3220"
120 },
121 {
122 "ImportPath": "github.com/jbenet/go-multihash",
Godeps/_workspace/src/github.com/h2so5/utp/.gitignore new
+26
@@ -0,0 +1,26 @@
1 +# Compiled Object files, Static and Dynamic libs (Shared Objects)
2 +*.o
3 +*.a
4 +*.so
5 +
6 +# Folders
7 +_obj
8 +_test
9 +
10 +# Architecture specific extensions/prefixes
11 +*.[568vq]
12 +[568vq].out
13 +
14 +*.cgo1.go
15 +*.cgo2.c
16 +_cgo_defun.c
17 +_cgo_gotypes.go
18 +_cgo_export.*
19 +
20 +_testmain.go
21 +
22 +*.exe
23 +*.test
24 +*.prof
25 +
26 +_ucat_test/libutp
Godeps/_workspace/src/github.com/h2so5/utp/.travis.yml new
+7
@@ -0,0 +1,7 @@
1 +language: go
2 +
3 +script:
4 + - GO_UTP_LOGGING=2 go test -v -bench .
5 + - go test -v -race
6 + - GO_UTP_LOGGING=2 go run benchmark/main.go -h
7 + - GO_UTP_LOGGING=2 cd _ucat_test; make test
Godeps/_workspace/src/github.com/h2so5/utp/LICENSE new
+21
@@ -0,0 +1,21 @@
1 +The MIT License (MIT)
2 +
3 +Copyright (c) 2014 Ron Hashimoto
4 +
5 +Permission is hereby granted, free of charge, to any person obtaining a copy
6 +of this software and associated documentation files (the "Software"), to deal
7 +in the Software without restriction, including without limitation the rights
8 +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
9 +copies of the Software, and to permit persons to whom the Software is
10 +furnished to do so, subject to the following conditions:
11 +
12 +The above copyright notice and this permission notice shall be included in all
13 +copies or substantial portions of the Software.
14 +
15 +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
16 +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
17 +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
18 +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
19 +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
20 +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
21 +SOFTWARE.
Godeps/_workspace/src/github.com/h2so5/utp/README.md new
+57
@@ -0,0 +1,57 @@
1 +utp
2 +===
3 +
4 +μTP (Micro Transport Protocol) implementation
5 +
6 +[![Build status](https://ci.appveyor.com/api/projects/status/j1be8y7p6nd2wqqw?svg=true)](https://ci.appveyor.com/project/h2so5/utp)
7 +[![Build Status](https://travis-ci.org/h2so5/utp.svg)](https://travis-ci.org/h2so5/utp)
8 +[![GoDoc](https://godoc.org/github.com/h2so5/utp?status.svg)](http://godoc.org/github.com/h2so5/utp)
9 +
10 +http://www.bittorrent.org/beps/bep_0029.html
11 +
12 +**warning: This is a buggy alpha version.**
13 +
14 +## Benchmark History
15 +
16 +[![Benchmark status](http://107.170.244.57:80/go-utp-bench.php)]()
17 +
18 +## Installation
19 +
20 +```
21 +go get github.com/h2so5/utp
22 +```
23 +
24 +## Example
25 +
26 +Echo server
27 +
28 +```go
29 +package main
30 +
31 +import (
32 + "time"
33 +
34 + "github.com/h2so5/utp"
35 +)
36 +
37 +func main() {
38 + ln, _ := utp.Listen("utp", ":11000")
39 + defer ln.Close()
40 +
41 + conn, _ := ln.AcceptUTP()
42 + conn.SetKeepAlive(time.Minute)
43 + defer conn.Close()
44 +
45 + for {
46 + var buf [1024]byte
47 + l, err := conn.Read(buf[:])
48 + if err != nil {
49 + break
50 + }
51 + _, err = conn.Write(buf[:l])
52 + if err != nil {
53 + break
54 + }
55 + }
56 +}
57 +```
Godeps/_workspace/src/github.com/h2so5/utp/addr.go new
+34
@@ -0,0 +1,34 @@
1 +package utp
2 +
3 +import "net"
4 +
5 +type UTPAddr struct {
6 + net.Addr
7 +}
8 +
9 +func (a UTPAddr) Network() string { return "utp" }
10 +
11 +func utp2udp(n string) (string, error) {
12 + switch n {
13 + case "utp":
14 + return "udp", nil
15 + case "utp4":
16 + return "udp4", nil
17 + case "utp6":
18 + return "udp6", nil
19 + default:
20 + return "", net.UnknownNetworkError(n)
21 + }
22 +}
23 +
24 +func ResolveUTPAddr(n, addr string) (*UTPAddr, error) {
25 + udpnet, err := utp2udp(n)
26 + if err != nil {
27 + return nil, err
28 + }
29 + udp, err := net.ResolveUDPAddr(udpnet, addr)
30 + if err != nil {
31 + return nil, err
32 + }
33 + return &UTPAddr{Addr: udp}, nil
34 +}
Godeps/_workspace/src/github.com/h2so5/utp/benchmark/main.go new
+276
@@ -0,0 +1,276 @@
1 +package main
2 +
3 +import (
4 + "bytes"
5 + "crypto/md5"
6 + "flag"
7 + "fmt"
8 + "io"
9 + "log"
10 + "math/rand"
11 + "sync"
12 + "time"
13 +
14 + "github.com/davecheney/profile"
15 + "github.com/dustin/go-humanize"
16 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/h2so5/utp"
17 +)
18 +
19 +type RandReader struct{}
20 +
21 +func (r RandReader) Read(p []byte) (n int, err error) {
22 + for i := range p {
23 + p[i] = byte(rand.Int())
24 + }
25 + return len(p), nil
26 +}
27 +
28 +type ByteCounter struct {
29 + n int64
30 + mutex sync.RWMutex
31 +}
32 +
33 +func (b *ByteCounter) Write(p []byte) (n int, err error) {
34 + b.mutex.Lock()
35 + defer b.mutex.Unlock()
36 + b.n += int64(len(p))
37 + return len(p), nil
38 +}
39 +
40 +func (b *ByteCounter) Length() int64 {
41 + b.mutex.RLock()
42 + defer b.mutex.RUnlock()
43 + return b.n
44 +}
45 +
46 +var h = flag.Bool("h", false, "Human readable")
47 +
48 +func main() {
49 + var l = flag.Int("c", 10485760, "Payload length (bytes)")
50 + var s = flag.Bool("s", false, "Stream mode(Low memory usage, but Slow)")
51 + flag.Parse()
52 +
53 + defer profile.Start(profile.CPUProfile).Stop()
54 +
55 + if *h {
56 + fmt.Printf("Payload: %s\n", humanize.IBytes(uint64(*l)))
57 + } else {
58 + fmt.Printf("Payload: %d\n", *l)
59 + }
60 +
61 + c2s := c2s(int64(*l), *s)
62 + n, p := humanize.ComputeSI(c2s)
63 + if *h {
64 + fmt.Printf("C2S: %f%sbps\n", n, p)
65 + } else {
66 + fmt.Printf("C2S: %f\n", c2s)
67 + }
68 +
69 + s2c := s2c(int64(*l), *s)
70 + n, p = humanize.ComputeSI(s2c)
71 + if *h {
72 + fmt.Printf("S2C: %f%sbps\n", n, p)
73 + } else {
74 + fmt.Printf("S2C: %f\n", s2c)
75 + }
76 +
77 + avg := (c2s + s2c) / 2.0
78 + n, p = humanize.ComputeSI(avg)
79 +
80 + if *h {
81 + fmt.Printf("AVG: %f%sbps\n", n, p)
82 + } else {
83 + fmt.Printf("AVG: %f\n", avg)
84 + }
85 +}
86 +
87 +func c2s(l int64, stream bool) float64 {
88 + ln, err := utp.Listen("utp", "127.0.0.1:0")
89 + if err != nil {
90 + log.Fatal(err)
91 + }
92 +
93 + raddr, err := utp.ResolveUTPAddr("utp", ln.Addr().String())
94 + if err != nil {
95 + log.Fatal(err)
96 + }
97 +
98 + c, err := utp.DialUTPTimeout("utp", nil, raddr, 1000*time.Millisecond)
99 + if err != nil {
100 + log.Fatal(err)
101 + }
102 + defer c.Close()
103 +
104 + if err != nil {
105 + log.Fatal(err)
106 + }
107 +
108 + s, err := ln.Accept()
109 + if err != nil {
110 + log.Fatal(err)
111 + }
112 + defer s.Close()
113 + ln.Close()
114 +
115 + rch := make(chan int)
116 +
117 + sendHash := md5.New()
118 + readHash := md5.New()
119 + counter := ByteCounter{}
120 +
121 + var bps float64
122 + if stream {
123 + go func() {
124 + defer c.Close()
125 + io.Copy(io.MultiWriter(c, sendHash, &counter), io.LimitReader(RandReader{}, l))
126 + }()
127 +
128 + go func() {
129 + io.Copy(readHash, s)
130 + close(rch)
131 + }()
132 +
133 + go func() {
134 + for {
135 + select {
136 + case <-time.After(time.Second):
137 + if *h {
138 + fmt.Printf("\r <--> %s ", humanize.IBytes(uint64(counter.Length())))
139 + } else {
140 + fmt.Printf("\r <--> %d ", counter.Length())
141 + }
142 + case <-rch:
143 + fmt.Printf("\r")
144 + return
145 + }
146 + }
147 + }()
148 +
149 + start := time.Now()
150 + <-rch
151 + bps = float64(l*8) / (float64(time.Now().Sub(start)) / float64(time.Second))
152 +
153 + } else {
154 + var sendBuf, readBuf bytes.Buffer
155 + io.Copy(io.MultiWriter(&sendBuf, sendHash), io.LimitReader(RandReader{}, l))
156 +
157 + go func() {
158 + defer c.Close()
159 + io.Copy(c, &sendBuf)
160 + }()
161 +
162 + go func() {
163 + io.Copy(&readBuf, s)
164 + rch <- 0
165 + }()
166 +
167 + start := time.Now()
168 + <-rch
169 + bps = float64(l*8) / (float64(time.Now().Sub(start)) / float64(time.Second))
170 +
171 + io.Copy(sendHash, &sendBuf)
172 + io.Copy(readHash, &readBuf)
173 + }
174 +
175 + if !bytes.Equal(sendHash.Sum(nil), readHash.Sum(nil)) {
176 + log.Fatal("Broken payload")
177 + }
178 +
179 + return bps
180 +}
181 +
182 +func s2c(l int64, stream bool) float64 {
183 + ln, err := utp.Listen("utp", "127.0.0.1:0")
184 + if err != nil {
185 + log.Fatal(err)
186 + }
187 +
188 + raddr, err := utp.ResolveUTPAddr("utp", ln.Addr().String())
189 + if err != nil {
190 + log.Fatal(err)
191 + }
192 +
193 + c, err := utp.DialUTPTimeout("utp", nil, raddr, 1000*time.Millisecond)
194 + if err != nil {
195 + log.Fatal(err)
196 + }
197 + defer c.Close()
198 +
199 + if err != nil {
200 + log.Fatal(err)
201 + }
202 +
203 + s, err := ln.Accept()
204 + if err != nil {
205 + log.Fatal(err)
206 + }
207 + defer s.Close()
208 + ln.Close()
209 +
210 + rch := make(chan int)
211 +
212 + sendHash := md5.New()
213 + readHash := md5.New()
214 + counter := ByteCounter{}
215 +
216 + var bps float64
217 +
218 + if stream {
219 + go func() {
220 + defer s.Close()
221 + io.Copy(io.MultiWriter(s, sendHash, &counter), io.LimitReader(RandReader{}, l))
222 + }()
223 +
224 + go func() {
225 + io.Copy(readHash, c)
226 + close(rch)
227 + }()
228 +
229 + go func() {
230 + for {
231 + select {
232 + case <-time.After(time.Second):
233 + if *h {
234 + fmt.Printf("\r <--> %s ", humanize.IBytes(uint64(counter.Length())))
235 + } else {
236 + fmt.Printf("\r <--> %d ", counter.Length())
237 + }
238 + case <-rch:
239 + fmt.Printf("\r")
240 + return
241 + }
242 + }
243 + }()
244 +
245 + start := time.Now()
246 + <-rch
247 + bps = float64(l*8) / (float64(time.Now().Sub(start)) / float64(time.Second))
248 +
249 + } else {
250 + var sendBuf, readBuf bytes.Buffer
251 + io.Copy(io.MultiWriter(&sendBuf, sendHash), io.LimitReader(RandReader{}, l))
252 +
253 + go func() {
254 + defer s.Close()
255 + io.Copy(s, &sendBuf)
256 + }()
257 +
258 + go func() {
259 + io.Copy(&readBuf, c)
260 + rch <- 0
261 + }()
262 +
263 + start := time.Now()
264 + <-rch
265 + bps = float64(l*8) / (float64(time.Now().Sub(start)) / float64(time.Second))
266 +
267 + io.Copy(sendHash, &sendBuf)
268 + io.Copy(readHash, &readBuf)
269 + }
270 +
271 + if !bytes.Equal(sendHash.Sum(nil), readHash.Sum(nil)) {
272 + log.Fatal("Broken payload")
273 + }
274 +
275 + return bps
276 +}
Godeps/_workspace/src/github.com/h2so5/utp/buffer.go new
+230
@@ -0,0 +1,230 @@
1 +package utp
2 +
3 +import (
4 + "errors"
5 + "math"
6 + "time"
7 +)
8 +
9 +type packetBuffer struct {
10 + root *packetBufferNode
11 + size int
12 + begin int
13 +}
14 +
15 +type packetBufferNode struct {
16 + p *packet
17 + next *packetBufferNode
18 + pushed time.Time
19 +}
20 +
21 +func newPacketBuffer(size, begin int) *packetBuffer {
22 + return &packetBuffer{
23 + size: size,
24 + begin: begin,
25 + }
26 +}
27 +
28 +func (b *packetBuffer) push(p *packet) error {
29 + if int(p.header.seq) > b.begin+b.size-1 {
30 + return errors.New("out of bounds")
31 + } else if int(p.header.seq) < b.begin {
32 + if int(p.header.seq)+math.MaxUint16 > b.begin+b.size-1 {
33 + return errors.New("out of bounds")
34 + }
35 + }
36 + if b.root == nil {
37 + b.root = &packetBufferNode{}
38 + }
39 + n := b.root
40 + i := b.begin
41 + for {
42 + if i == int(p.header.seq) {
43 + n.p = p
44 + n.pushed = time.Now()
45 + return nil
46 + } else if n.next == nil {
47 + n.next = &packetBufferNode{}
48 + }
49 + n = n.next
50 + i = (i + 1) % (math.MaxUint16 + 1)
51 + }
52 + return nil
53 +}
54 +
55 +func (b *packetBuffer) fetch(id uint16) *packet {
56 + for p := b.root; p != nil; p = p.next {
57 + if p.p != nil {
58 + if p.p.header.seq < id {
59 + p.p = nil
60 + } else if p.p.header.seq == id {
61 + r := p.p
62 + p.p = nil
63 + return r
64 + }
65 + }
66 + }
67 + return nil
68 +}
69 +
70 +func (b *packetBuffer) compact() {
71 + for b.root != nil && b.root.p == nil {
72 + b.root = b.root.next
73 + b.begin = (b.begin + 1) % (math.MaxUint16 + 1)
74 + }
75 +}
76 +
77 +func (b *packetBuffer) first() *packet {
78 + if b.root == nil || b.root.p == nil {
79 + return nil
80 + }
81 + return b.root.p
82 +}
83 +
84 +func (b *packetBuffer) frontPushedTime() (time.Time, error) {
85 + if b.root == nil || b.root.p == nil {
86 + return time.Time{}, errors.New("no first packet")
87 + }
88 + return b.root.pushed, nil
89 +}
90 +
91 +func (b *packetBuffer) fetchSequence() []*packet {
92 + var a []*packet
93 + for ; b.root != nil && b.root.p != nil; b.root = b.root.next {
94 + a = append(a, b.root.p)
95 + b.begin = (b.begin + 1) % (math.MaxUint16 + 1)
96 + }
97 + return a
98 +}
99 +
100 +func (b *packetBuffer) sequence() []*packet {
101 + var a []*packet
102 + n := b.root
103 + for ; n != nil && n.p != nil; n = n.next {
104 + a = append(a, n.p)
105 + }
106 + return a
107 +}
108 +
109 +func (b *packetBuffer) space() int {
110 + s := b.size
111 + for p := b.root; p != nil; p = p.next {
112 + s--
113 + }
114 + return s
115 +}
116 +
117 +func (b *packetBuffer) empty() bool {
118 + return b.root == nil
119 +}
120 +
121 +// test use only
122 +func (b *packetBuffer) all() []*packet {
123 + var a []*packet
124 + for p := b.root; p != nil; p = p.next {
125 + if p.p != nil {
126 + a = append(a, p.p)
127 + }
128 + }
129 + return a
130 +}
131 +
132 +func (b *packetBuffer) generateSelectiveACK() []byte {
133 + if b.empty() {
134 + return nil
135 + }
136 +
137 + var ack []byte
138 + var bit uint
139 + var octet byte
140 + for p := b.root.next; p != nil; p = p.next {
141 + if p.p != nil {
142 + octet |= (1 << bit)
143 + }
144 + bit++
145 + if bit == 8 {
146 + ack = append(ack, octet)
147 + bit = 0
148 + octet = 0
149 + }
150 + }
151 +
152 + if bit != 0 {
153 + ack = append(ack, octet)
154 + }
155 +
156 + for len(ack) > 0 && ack[len(ack)-1] == 0 {
157 + ack = ack[:len(ack)-1]
158 + }
159 +
160 + return ack
161 +}
162 +
163 +func (b *packetBuffer) processSelectiveACK(ack []byte) {
164 + if b.empty() {
165 + return
166 + }
167 +
168 + p := b.root.next
169 + if p == nil {
170 + return
171 + }
172 +
173 + for _, a := range ack {
174 + for i := 0; i < 8; i++ {
175 + acked := (a & 1) != 0
176 + a >>= 1
177 + if acked {
178 + p.p = nil
179 + }
180 + p = p.next
181 + if p == nil {
182 + return
183 + }
184 + }
185 + }
186 +}
187 +
188 +type timedBuffer struct {
189 + d time.Duration
190 + root *timedBufferNode
191 +}
192 +
193 +type timedBufferNode struct {
194 + val float64
195 + next *timedBufferNode
196 + pushed time.Time
197 +}
198 +
199 +func (b *timedBuffer) push(val float64) {
200 + var before *timedBufferNode
201 + for n := b.root; n != nil; n = n.next {
202 + if time.Now().Sub(n.pushed) >= b.d {
203 + if before != nil {
204 + before.next = nil
205 + } else {
206 + b.root = nil
207 + }
208 + break
209 + }
210 + before = n
211 + }
212 + b.root = &timedBufferNode{
213 + val: val,
214 + next: b.root,
215 + pushed: time.Now(),
216 + }
217 +}
218 +
219 +func (b *timedBuffer) min() float64 {
220 + if b.root == nil {
221 + return 0
222 + }
223 + min := b.root.val
224 + for n := b.root; n != nil; n = n.next {
225 + if min > n.val {
226 + min = n.val
227 + }
228 + }
229 + return min
230 +}
Godeps/_workspace/src/github.com/h2so5/utp/conn.go new
+761
@@ -0,0 +1,761 @@
1 +package utp
2 +
3 +import (
4 + "bytes"
5 + "errors"
6 + "io"
7 + "math"
8 + "math/rand"
9 + "net"
10 + "syscall"
11 + "time"
12 +)
13 +
14 +type UTPConn struct {
15 + conn net.PacketConn
16 + raddr net.Addr
17 + rid, sid, seq, ack, lastAck uint16
18 + rtt, rttVar, minRtt, rto, dupAck int64
19 + diff, maxWindow uint32
20 + rdeadline, wdeadline time.Time
21 +
22 + state state
23 + lastTimedOut time.Time
24 +
25 + outch chan *outgoingPacket
26 + outchch chan int
27 + sendch chan *outgoingPacket
28 + sendchch chan int
29 + recvch chan *packet
30 + recvchch chan int
31 + readch chan []byte
32 + readchch chan int
33 + winch chan uint32
34 + quitch chan int
35 + activech chan int
36 + connch chan error
37 + finch chan int
38 + closech chan<- uint16
39 + eofid uint16
40 + keepalivech chan time.Duration
41 +
42 + readbuf bytes.Buffer
43 + recvbuf *packetBuffer
44 + sendbuf *packetBuffer
45 +
46 + stat statistics
47 +}
48 +
49 +type statistics struct {
50 + sentPackets int
51 + resentPackets int
52 + receivedPackets int
53 + receivedDuplicatedACKs int
54 + packetTimedOuts int
55 + sentSelectiveACKs int
56 + receivedSelectiveACKs int
57 +
58 + rtoSum int
59 + rtoCount int
60 +}
61 +
62 +func dial(n string, laddr, raddr *UTPAddr, timeout time.Duration) (*UTPConn, error) {
63 + udpnet, err := utp2udp(n)
64 + if err != nil {
65 + return nil, err
66 + }
67 +
68 + // TODO extract
69 + if laddr == nil {
70 + addr, err := net.ResolveUDPAddr(udpnet, ":0")
71 + if err != nil {
72 + return nil, err
73 + }
74 + laddr = &UTPAddr{Addr: addr}
75 + }
76 +
77 + conn, err := net.ListenPacket(udpnet, laddr.Addr.String())
78 + if err != nil {
79 + return nil, err
80 + }
81 +
82 + id := uint16(rand.Intn(math.MaxUint16))
83 +
84 + c := newUTPConn()
85 + c.conn = conn
86 + c.raddr = raddr.Addr
87 + c.rid = id
88 + c.sid = id + 1
89 + c.seq = 1
90 + c.state = state_syn_sent
91 + c.sendbuf = newPacketBuffer(window_size, 1)
92 +
93 + go c.recv()
94 + go c.loop()
95 +
96 + select {
97 + case c.sendch <- &outgoingPacket{st_syn, nil, nil}:
98 + case <-c.sendchch:
99 + return nil, errors.New("use of closed network connection")
100 + }
101 +
102 + var t <-chan time.Time
103 + if timeout != 0 {
104 + t = time.After(timeout)
105 + }
106 +
107 + select {
108 + case err := <-c.connch:
109 + if err != nil {
110 + c.closed()
111 + return nil, err
112 + }
113 + ulog.Printf(1, "Conn(%v): Connected", c.LocalAddr())
114 + return c, nil
115 + case <-t:
116 + c.quitch <- 0
117 + return nil, &timeoutError{}
118 + }
119 +}
120 +
121 +func newUTPConn() *UTPConn {
122 + rto := 60
123 +
124 + return &UTPConn{
125 + minRtt: math.MaxInt64,
126 + maxWindow: mtu,
127 + rto: int64(rto),
128 +
129 + outch: make(chan *outgoingPacket, 1),
130 + outchch: make(chan int),
131 + sendch: make(chan *outgoingPacket, 1),
132 + sendchch: make(chan int),
133 + recvch: make(chan *packet, 2),
134 + recvchch: make(chan int),
135 + winch: make(chan uint32, 1),
136 + quitch: make(chan int),
137 + activech: make(chan int),
138 + readch: make(chan []byte, 1),
139 + readchch: make(chan int),
140 + connch: make(chan error, 1),
141 + finch: make(chan int, 1),
142 +
143 + keepalivech: make(chan time.Duration),
144 +
145 + stat: statistics{
146 + rtoSum: rto,
147 + rtoCount: 1,
148 + },
149 + }
150 +}
151 +
152 +func (c *UTPConn) ok() bool { return c != nil && c.conn != nil }
153 +
154 +func (c *UTPConn) Close() error {
155 + if !c.ok() {
156 + return syscall.EINVAL
157 + }
158 +
159 + select {
160 + case <-c.activech:
161 + default:
162 + c.quitch <- 0
163 + ulog.Printf(2, "Conn(%v): Wait for close", c.LocalAddr())
164 + <-c.finch
165 + }
166 +
167 + return nil
168 +}
169 +
170 +func (c *UTPConn) LocalAddr() net.Addr {
171 + return &UTPAddr{Addr: c.conn.LocalAddr()}
172 +}
173 +
174 +func (c *UTPConn) RemoteAddr() net.Addr {
175 + return &UTPAddr{Addr: c.raddr}
176 +}
177 +
178 +func (c *UTPConn) Read(b []byte) (int, error) {
179 + if !c.ok() {
180 + return 0, syscall.EINVAL
181 + }
182 +
183 + if c.readbuf.Len() == 0 {
184 + var timeout <-chan time.Time
185 + if !c.rdeadline.IsZero() {
186 + timeout = time.After(c.rdeadline.Sub(time.Now()))
187 + }
188 +
189 + select {
190 + case b := <-c.readch:
191 + if b == nil {
192 + return 0, io.EOF
193 + }
194 + _, err := c.readbuf.Write(b)
195 + if err != nil {
196 + return 0, err
197 + }
198 + case <-c.readchch:
199 + loop:
200 + for {
201 + select {
202 + case b := <-c.readch:
203 + _, err := c.readbuf.Write(b)
204 + if err != nil {
205 + return 0, err
206 + }
207 + default:
208 + break loop
209 + }
210 + }
211 + if c.readbuf.Len() == 0 {
212 + return 0, io.EOF
213 + }
214 + case <-timeout:
215 + return 0, &timeoutError{}
216 + }
217 + }
218 + return c.readbuf.Read(b)
219 +}
220 +
221 +func (c *UTPConn) Write(b []byte) (int, error) {
222 + if !c.ok() {
223 + return 0, syscall.EINVAL
224 + }
225 +
226 + var wrote uint64
227 + for {
228 + l := uint64(len(b)) - wrote
229 + if l > mss {
230 + l = mss
231 + }
232 + select {
233 + case c.outch <- &outgoingPacket{st_data, nil, b[wrote : wrote+l]}:
234 + case <-c.outchch:
235 + return 0, errors.New("use of closed network connection")
236 + }
237 +
238 + wrote += l
239 + ulog.Printf(4, "Conn(%v): Write %d/%d bytes", c.LocalAddr(), wrote, len(b))
240 + if l < mss {
241 + break
242 + }
243 + }
244 +
245 + return len(b), nil
246 +}
247 +
248 +func (c *UTPConn) SetDeadline(t time.Time) error {
249 + if !c.ok() {
250 + return syscall.EINVAL
251 + }
252 + if err := c.SetReadDeadline(t); err != nil {
253 + return err
254 + }
255 + if err := c.SetWriteDeadline(t); err != nil {
256 + return err
257 + }
258 + return nil
259 +}
260 +
261 +func (c *UTPConn) SetReadDeadline(t time.Time) error {
262 + if !c.ok() {
263 + return syscall.EINVAL
264 + }
265 + c.rdeadline = t
266 + return nil
267 +}
268 +
269 +func (c *UTPConn) SetWriteDeadline(t time.Time) error {
270 + if !c.ok() {
271 + return syscall.EINVAL
272 + }
273 + c.wdeadline = t
274 + return nil
275 +}
276 +
277 +func (c *UTPConn) SetKeepAlive(d time.Duration) error {
278 + if !c.ok() {
279 + return syscall.EINVAL
280 + }
281 + select {
282 + case <-c.activech:
283 + default:
284 + c.keepalivech <- d
285 + }
286 + return nil
287 +}
288 +
289 +func readPacket(data []byte) (*packet, error) {
290 + p := globalPool.get()
291 + err := p.UnmarshalBinary(data)
292 + if err != nil {
293 + return nil, err
294 + }
295 + if p.header.ver != version {
296 + return nil, errors.New("unsupported header version")
297 + }
298 + return p, nil
299 +}
300 +
301 +func (c *UTPConn) recv() {
302 + for {
303 + var buf [mtu]byte
304 + len, addr, err := c.conn.ReadFrom(buf[:])
305 + if err != nil {
306 + return
307 + }
308 + if addr.String() != c.raddr.String() {
309 + continue
310 + }
311 + p, err := readPacket(buf[:len])
312 + if err == nil {
313 + select {
314 + case c.recvch <- p:
315 + case <-c.recvchch:
316 + return
317 + }
318 + }
319 + }
320 +}
321 +
322 +func (c *UTPConn) loop() {
323 + var recvExit, sendExit bool
324 + var lastReceived time.Time
325 + var keepalive <-chan time.Time
326 +
327 + go func() {
328 + var window uint32 = window_size * mtu
329 + for {
330 + if window >= mtu {
331 + select {
332 + case b := <-c.outch:
333 + select {
334 + case c.sendch <- b:
335 + window -= mtu
336 + case <-c.sendchch:
337 + return
338 + }
339 + case <-c.outchch:
340 + return
341 + case w := <-c.winch:
342 + window = w
343 + }
344 + } else {
345 + window = <-c.winch
346 + }
347 + }
348 + }()
349 +
350 + for {
351 + select {
352 + case <-c.sendchch:
353 + sendExit = true
354 + default:
355 + }
356 + select {
357 + case <-c.recvchch:
358 + recvExit = true
359 + default:
360 + }
361 + select {
362 + case p := <-c.recvch:
363 + ack := c.processPacket(p)
364 + lastReceived = time.Now()
365 + if ack {
366 + out := &outgoingPacket{st_state, nil, nil}
367 + selack := c.sendbuf.generateSelectiveACK()
368 + if len(selack) > 0 {
369 + out.ext = []extension{
370 + extension{
371 + typ: ext_selective_ack,
372 + payload: selack,
373 + },
374 + }
375 + c.stat.sentSelectiveACKs++
376 + }
377 + c.sendPacket(out)
378 + }
379 +
380 + case b := <-c.sendch:
381 + c.sendPacket(b)
382 +
383 + case <-time.After(time.Duration(c.rto) * time.Millisecond):
384 + if !c.state.active && time.Now().Sub(lastReceived) > reset_timeout {
385 + ulog.Printf(2, "Conn(%v): Connection timed out", c.LocalAddr())
386 + c.sendPacket(&outgoingPacket{st_reset, nil, nil})
387 + c.close()
388 + } else {
389 + t, err := c.sendbuf.frontPushedTime()
390 + if err == nil && c.lastTimedOut != t && time.Now().Sub(t) > time.Duration(c.rto)*time.Millisecond {
391 + c.lastTimedOut = t
392 + c.stat.packetTimedOuts++
393 + c.maxWindow /= 2
394 + if c.maxWindow < mtu {
395 + c.maxWindow = mtu
396 + }
397 + for _, p := range c.sendbuf.sequence() {
398 + c.resendPacket(p)
399 + }
400 + }
401 + }
402 + case d := <-c.keepalivech:
403 + if d <= 0 {
404 + keepalive = nil
405 + } else {
406 + keepalive = time.Tick(d)
407 + }
408 + case <-keepalive:
409 + ulog.Printf(2, "Conn(%v): Send keepalive", c.LocalAddr())
410 + c.sendPacket(&outgoingPacket{st_state, nil, nil})
411 +
412 + case <-c.quitch:
413 + if c.state.exit != nil {
414 + c.state.exit(c)
415 + }
416 + }
417 + if recvExit && sendExit {
418 + return
419 + }
420 + }
421 +}
422 +
423 +func (c *UTPConn) sendPacket(b *outgoingPacket) {
424 + p := c.makePacket(b)
425 + bin, err := p.MarshalBinary()
426 + if err == nil {
427 + ulog.Printf(3, "SEND %v -> %v: %v", c.conn.LocalAddr(), c.raddr, p.String())
428 + c.stat.sentPackets++
429 + _, err = c.conn.WriteTo(bin, c.raddr)
430 + if err != nil {
431 + return
432 + }
433 + if b.typ != st_state {
434 + c.sendbuf.push(p)
435 + } else {
436 + globalPool.put(p)
437 + }
438 + }
439 +}
440 +
441 +func (c *UTPConn) resendPacket(p *packet) {
442 + bin, err := p.MarshalBinary()
443 + if err == nil {
444 + ulog.Printf(3, "RESEND %v -> %v: %v", c.conn.LocalAddr(), c.raddr, p.String())
445 + c.stat.resentPackets++
446 + _, err = c.conn.WriteTo(bin, c.raddr)
447 + if err != nil {
448 + return
449 + }
450 + }
451 +}
452 +
453 +func currentMicrosecond() uint32 {
454 + return uint32(time.Now().Nanosecond() / 1000)
455 +}
456 +
457 +func (c *UTPConn) processPacket(p *packet) bool {
458 + var ack bool
459 +
460 + if p.header.t == 0 {
461 + c.diff = 0
462 + } else {
463 + t := currentMicrosecond()
464 + if t > p.header.t {
465 + c.diff = t - p.header.t
466 + if c.minRtt > int64(c.diff) {
467 + c.minRtt = int64(c.diff)
468 + }
469 + }
470 + }
471 +
472 + ulog.Printf(3, "RECV %v -> %v: %v", c.raddr, c.conn.LocalAddr(), p.String())
473 + c.stat.receivedPackets++
474 +
475 + if p.header.typ == st_state {
476 +
477 + f := c.sendbuf.first()
478 + if f != nil && p.header.ack == f.header.seq {
479 + for _, e := range p.ext {
480 + if e.typ == ext_selective_ack {
481 + ulog.Printf(3, "Conn(%v): Receive Selective ACK", c.LocalAddr())
482 + c.stat.receivedSelectiveACKs++
483 + c.sendbuf.processSelectiveACK(e.payload)
484 + }
485 + }
486 + }
487 +
488 + s := c.sendbuf.fetch(p.header.ack)
489 + if s != nil {
490 + current := currentMicrosecond()
491 + if current > s.header.t {
492 + e := int64(current-s.header.t) / 1000
493 + if c.rtt == 0 {
494 + c.rtt = e
495 + c.rttVar = e / 2
496 + } else {
497 + d := c.rtt - e
498 + if d < 0 {
499 + d = -d
500 + }
501 + c.rttVar += (d - c.rttVar) / 4
502 + c.rtt = c.rtt - c.rtt/8 + e/8
503 + }
504 + c.rto = c.rtt + c.rttVar*4
505 + if c.rto < 60 {
506 + c.rto = 60
507 + } else if c.rto > 1000 {
508 + c.rto = 1000
509 + }
510 + c.stat.rtoSum += int(c.rto)
511 + c.stat.rtoCount++
512 + }
513 +
514 + if c.diff != 0 {
515 + ourDelay := float64(c.diff)
516 + offTarget := 100000.0 - ourDelay
517 + windowFactor := float64(mtu) / float64(c.maxWindow)
518 + delayFactor := offTarget / 100000.0
519 + gain := 3000.0 * delayFactor * windowFactor
520 + c.maxWindow = uint32(int(c.maxWindow) + int(gain))
521 + if c.maxWindow < mtu {
522 + c.maxWindow = mtu
523 + }
524 + ulog.Printf(4, "Conn(%v): Update maxWindow: %d", c.LocalAddr(), c.maxWindow)
525 + }
526 + globalPool.put(s)
527 + }
528 + c.sendbuf.compact()
529 + if c.lastAck == p.header.ack {
530 + c.dupAck++
531 + if c.dupAck >= 2 {
532 + ulog.Printf(3, "Conn(%v): Receive 3 duplicated acks: %d", c.LocalAddr(), p.header.ack)
533 + c.stat.receivedDuplicatedACKs++
534 + p := c.sendbuf.first()
535 + if p != nil {
536 + c.maxWindow /= 2
537 + if c.maxWindow < mtu {
538 + c.maxWindow = mtu
539 + }
540 + ulog.Printf(4, "Conn(%v): Update maxWindow: %d", c.LocalAddr(), c.maxWindow)
541 + c.resendPacket(p)
542 + }
543 + c.dupAck = 0
544 + }
545 + } else {
546 + c.dupAck = 0
547 + }
548 + c.lastAck = p.header.ack
549 + if p.header.ack == c.seq-1 {
550 + wnd := p.header.wnd
551 + if wnd > c.maxWindow {
552 + wnd = c.maxWindow
553 + }
554 + ulog.Printf(4, "Conn(%v): Reset window: %d", c.LocalAddr(), wnd)
555 + go func() {
556 + c.winch <- wnd
557 + }()
558 + }
559 + if c.state.state != nil {
560 + c.state.state(c, p)
561 + }
562 + globalPool.put(p)
563 + } else if p.header.typ == st_reset {
564 + globalPool.put(p)
565 + c.close()
566 + } else {
567 + if c.recvbuf == nil {
568 + return false
569 + }
570 + ack = true
571 + c.recvbuf.push(p)
572 + for _, s := range c.recvbuf.fetchSequence() {
573 + c.ack = s.header.seq
574 + switch s.header.typ {
575 + case st_data:
576 + if c.state.data != nil {
577 + c.state.data(c, s)
578 + }
579 + case st_fin:
580 + if c.state.fin != nil {
581 + c.state.fin(c, s)
582 + }
583 + case st_state:
584 + if c.state.state != nil {
585 + c.state.state(c, s)
586 + }
587 + }
588 + globalPool.put(s)
589 + }
590 + }
591 + return ack
592 +}
593 +
594 +func (c *UTPConn) makePacket(b *outgoingPacket) *packet {
595 + wnd := window_size * mtu
596 + if c.recvbuf != nil {
597 + wnd = c.recvbuf.space() * mtu
598 + }
599 + id := c.sid
600 + if b.typ == st_syn {
601 + id = c.rid
602 + }
603 + p := globalPool.get()
604 + p.header.typ = b.typ
605 + p.header.ver = version
606 + p.header.id = id
607 + p.header.t = currentMicrosecond()
608 + p.header.diff = c.diff
609 + p.header.wnd = uint32(wnd)
610 + p.header.seq = c.seq
611 + p.header.ack = c.ack
612 + if b.typ == st_fin {
613 + c.eofid = c.seq
614 + }
615 + if !(b.typ == st_state && len(b.payload) == 0) {
616 + c.seq++
617 + }
618 + p.payload = p.payload[:len(b.payload)]
619 + copy(p.payload, b.payload)
620 + return p
621 +}
622 +
623 +func (c *UTPConn) close() {
624 + if !c.state.closed {
625 + close(c.outchch)
626 + close(c.readchch)
627 + close(c.sendchch)
628 + close(c.recvchch)
629 + close(c.activech)
630 + close(c.finch)
631 + c.closed()
632 +
633 + // Accepted connection
634 + if c.closech != nil {
635 + c.closech <- c.sid
636 + } else {
637 + c.conn.Close()
638 + }
639 +
640 + ulog.Printf(1, "Conn(%v): Closed", c.LocalAddr())
641 + ulog.Printf(1, "Conn(%v): * SentPackets: %d", c.LocalAddr(), c.stat.sentPackets)
642 + ulog.Printf(1, "Conn(%v): * ResentPackets: %d", c.LocalAddr(), c.stat.resentPackets)
643 + ulog.Printf(1, "Conn(%v): * ReceivedPackets: %d", c.LocalAddr(), c.stat.receivedPackets)
644 + ulog.Printf(1, "Conn(%v): * ReceivedDuplicatedACKs: %d", c.LocalAddr(), c.stat.receivedDuplicatedACKs)
645 + ulog.Printf(1, "Conn(%v): * PacketTimedOuts: %d", c.LocalAddr(), c.stat.packetTimedOuts)
646 + ulog.Printf(1, "Conn(%v): * SentSelectiveACKs: %d", c.LocalAddr(), c.stat.sentSelectiveACKs)
647 + ulog.Printf(1, "Conn(%v): * ReceivedSelectiveACKs: %d", c.LocalAddr(), c.stat.receivedSelectiveACKs)
648 + ulog.Printf(1, "Conn(%v): * AverageRTO: %d", c.LocalAddr(), c.stat.rtoSum/c.stat.rtoCount)
649 + }
650 +}
651 +
652 +func (c *UTPConn) closed() {
653 + ulog.Printf(2, "Conn(%v): Change state: CLOSED", c.LocalAddr())
654 + c.state = state_closed
655 +}
656 +
657 +func (c *UTPConn) closing() {
658 + ulog.Printf(2, "Conn(%v): Change state: CLOSING", c.LocalAddr())
659 + c.state = state_closing
660 +}
661 +
662 +func (c *UTPConn) syn_sent() {
663 + ulog.Printf(2, "Conn(%v): Change state: SYN_SENT", c.LocalAddr())
664 + c.state = state_syn_sent
665 +}
666 +
667 +func (c *UTPConn) connected() {
668 + ulog.Printf(2, "Conn(%v): Change state: CONNECTED", c.LocalAddr())
669 + c.state = state_connected
670 +}
671 +
672 +func (c *UTPConn) fin_sent() {
673 + ulog.Printf(2, "Conn(%v): Change state: FIN_SENT", c.LocalAddr())
674 + c.state = state_fin_sent
675 +}
676 +
677 +type state struct {
678 + data func(c *UTPConn, p *packet)
679 + fin func(c *UTPConn, p *packet)
680 + state func(c *UTPConn, p *packet)
681 + exit func(c *UTPConn)
682 + active bool
683 + closed bool
684 +}
685 +
686 +var state_closed state = state{
687 + closed: true,
688 +}
689 +
690 +var state_closing state = state{
691 + data: func(c *UTPConn, p *packet) {
692 + select {
693 + case c.readch <- append([]byte(nil), p.payload...):
694 + case <-c.readchch:
695 + }
696 + if c.recvbuf.empty() && c.sendbuf.empty() {
697 + c.close()
698 + }
699 + },
700 + state: func(c *UTPConn, p *packet) {
701 + if c.recvbuf.empty() && c.sendbuf.empty() {
702 + c.close()
703 + }
704 + },
705 +}
706 +
707 +var state_syn_sent state = state{
708 + state: func(c *UTPConn, p *packet) {
709 + c.recvbuf = newPacketBuffer(window_size, int(p.header.seq))
710 + c.connected()
711 + c.connch <- nil
712 + },
713 + exit: func(c *UTPConn) {
714 + go func() {
715 + select {
716 + case c.outch <- &outgoingPacket{st_fin, nil, nil}:
717 + case <-c.outchch:
718 + }
719 + }()
720 + c.fin_sent()
721 + },
722 + active: true,
723 +}
724 +
725 +var state_connected state = state{
726 + data: func(c *UTPConn, p *packet) {
727 + select {
728 + case c.readch <- append([]byte(nil), p.payload...):
729 + case <-c.readchch:
730 + }
731 + },
732 + fin: func(c *UTPConn, p *packet) {
733 + if c.recvbuf.empty() && c.sendbuf.empty() {
734 + c.close()
735 + } else {
736 + c.closing()
737 + }
738 + },
739 + exit: func(c *UTPConn) {
740 + go func() {
741 + select {
742 + case c.outch <- &outgoingPacket{st_fin, nil, nil}:
743 + case <-c.outchch:
744 + }
745 + }()
746 + c.fin_sent()
747 + },
748 + active: true,
749 +}
750 +
751 +var state_fin_sent state = state{
752 + state: func(c *UTPConn, p *packet) {
753 + if p.header.ack == c.eofid {
754 + if c.recvbuf.empty() && c.sendbuf.empty() {
755 + c.close()
756 + } else {
757 + c.closing()
758 + }
759 + }
760 + },
761 +}
Godeps/_workspace/src/github.com/h2so5/utp/dial.go new
+68
@@ -0,0 +1,68 @@
1 +package utp
2 +
3 +import (
4 + "errors"
5 + "net"
6 + "time"
7 +)
8 +
9 +func Dial(n, addr string) (*UTPConn, error) {
10 + raddr, err := ResolveUTPAddr(n, addr)
11 + if err != nil {
12 + return nil, err
13 + }
14 + return DialUTP(n, nil, raddr)
15 +}
16 +
17 +func DialUTP(n string, laddr, raddr *UTPAddr) (*UTPConn, error) {
18 + return dial(n, laddr, raddr, 0)
19 +}
20 +
21 +func DialUTPTimeout(n string, laddr, raddr *UTPAddr, timeout time.Duration) (*UTPConn, error) {
22 + return dial(n, laddr, raddr, timeout)
23 +}
24 +
25 +// A Dialer contains options for connecting to an address.
26 +//
27 +// The zero value for each field is equivalent to dialing without
28 +// that option. Dialing with the zero value of Dialer is therefore
29 +// equivalent to just calling the Dial function.
30 +type Dialer struct {
31 + // Timeout is the maximum amount of time a dial will wait for
32 + // a connect to complete. If Deadline is also set, it may fail
33 + // earlier.
34 + //
35 + // The default is no timeout.
36 + //
37 + // With or without a timeout, the operating system may impose
38 + // its own earlier timeout. For instance, TCP timeouts are
39 + // often around 3 minutes.
40 + Timeout time.Duration
41 +
42 + // LocalAddr is the local address to use when dialing an
43 + // address. The address must be of a compatible type for the
44 + // network being dialed.
45 + // If nil, a local address is automatically chosen.
46 + LocalAddr net.Addr
47 +}
48 +
49 +// Dial connects to the address on the named network.
50 +//
51 +// See func Dial for a description of the network and address parameters.
52 +func (d *Dialer) Dial(n, addr string) (*UTPConn, error) {
53 + raddr, err := ResolveUTPAddr(n, addr)
54 + if err != nil {
55 + return nil, err
56 + }
57 +
58 + var laddr *UTPAddr
59 + if d.LocalAddr != nil {
60 + var ok bool
61 + laddr, ok = d.LocalAddr.(*UTPAddr)
62 + if !ok {
63 + return nil, errors.New("Dialer.LocalAddr is not a UTPAddr")
64 + }
65 + }
66 +
67 + return DialUTPTimeout(n, laddr, raddr, d.Timeout)
68 +}
Godeps/_workspace/src/github.com/h2so5/utp/listener.go new
+329
@@ -0,0 +1,329 @@
1 +package utp
2 +
3 +import (
4 + "errors"
5 + "math"
6 + "math/rand"
7 + "net"
8 + "syscall"
9 + "time"
10 +)
11 +
12 +type UTPListener struct {
13 + // RawConn represents an out-of-band connection.
14 + // This allows a single socket to handle multiple protocols.
15 + RawConn net.PacketConn
16 +
17 + conn net.PacketConn
18 + conns map[uint16]*UTPConn
19 + accept chan (*UTPConn)
20 + err chan (error)
21 + lasterr error
22 + deadline time.Time
23 + closech chan int
24 + connch chan uint16
25 + closed bool
26 +}
27 +
28 +func Listen(n, laddr string) (*UTPListener, error) {
29 + addr, err := ResolveUTPAddr(n, laddr)
30 + if err != nil {
31 + return nil, err
32 + }
33 + return ListenUTP(n, addr)
34 +}
35 +
36 +func ListenUTP(n string, laddr *UTPAddr) (*UTPListener, error) {
37 + udpnet, err := utp2udp(n)
38 + if err != nil {
39 + return nil, err
40 + }
41 + conn, err := listenPacket(udpnet, laddr.Addr.String())
42 + if err != nil {
43 + return nil, err
44 + }
45 +
46 + l := UTPListener{
47 + RawConn: newRawConn(conn),
48 + conn: conn,
49 + conns: make(map[uint16]*UTPConn),
50 + accept: make(chan (*UTPConn), 10),
51 + err: make(chan (error), 1),
52 + closech: make(chan int),
53 + connch: make(chan uint16),
54 + lasterr: nil,
55 + }
56 +
57 + l.listen()
58 + return &l, nil
59 +}
60 +
61 +type incoming struct {
62 + p *packet
63 + addr net.Addr
64 +}
65 +
66 +func (l *UTPListener) listen() {
67 + inch := make(chan incoming)
68 + raw := l.RawConn.(*rawConn)
69 +
70 + // reads udp packets
71 + go func() {
72 + for {
73 + var buf [mtu]byte
74 + len, addr, err := l.conn.ReadFrom(buf[:])
75 + if err != nil {
76 + l.err <- err
77 + return
78 + }
79 + p, err := readPacket(buf[:len])
80 + if err == nil {
81 + inch <- incoming{p, addr}
82 + } else {
83 + select {
84 + case <-raw.closed:
85 + default:
86 + i := rawIncoming{b: buf[:len], addr: addr}
87 + select {
88 + case raw.in <- i:
89 + default:
90 + // discard the oldest packet
91 + <-raw.in
92 + raw.in <- i
93 + }
94 + }
95 + }
96 + }
97 + }()
98 +
99 + go func() {
100 + for {
101 + select {
102 + case i := <-inch:
103 + l.processPacket(i.p, i.addr)
104 + case <-l.closech:
105 + ulog.Printf(2, "Listener(%v): Stop listening", l.conn.LocalAddr())
106 + close(l.accept)
107 + l.closed = true
108 + case id := <-l.connch:
109 + if _, ok := l.conns[id]; !ok {
110 + delete(l.conns, id+1)
111 + ulog.Printf(2, "Listener(%v): Connection closed #%d (alive: %d)", l.conn.LocalAddr(), id, len(l.conns))
112 + if l.closed && len(l.conns) == 0 {
113 + ulog.Printf(2, "Listener(%v): All accepted connections are closed", l.conn.LocalAddr())
114 + l.conn.Close()
115 + ulog.Printf(1, "Listener(%v): Closed", l.conn.LocalAddr())
116 + return
117 + }
118 + }
119 + }
120 + }
121 + }()
122 +
123 + ulog.Printf(1, "Listener(%v): Start listening", l.conn.LocalAddr())
124 +}
125 +
126 +func listenPacket(n, addr string) (net.PacketConn, error) {
127 + if n == "mem" {
128 + return nil, errors.New("TODO implement in-memory packet connection")
129 + }
130 + return net.ListenPacket(n, addr)
131 +}
132 +
133 +func (l *UTPListener) processPacket(p *packet, addr net.Addr) {
134 + switch p.header.typ {
135 + case st_data, st_fin, st_state, st_reset:
136 + if c, ok := l.conns[p.header.id]; ok {
137 + select {
138 + case c.recvch <- p:
139 + case <-c.recvchch:
140 + }
141 + }
142 + case st_syn:
143 + if l.closed {
144 + return
145 + }
146 + sid := p.header.id + 1
147 + if _, ok := l.conns[p.header.id]; !ok {
148 + seq := rand.Intn(math.MaxUint16)
149 +
150 + c := newUTPConn()
151 + c.conn = l.conn
152 + c.raddr = addr
153 + c.rid = p.header.id + 1
154 + c.sid = p.header.id
155 + c.seq = uint16(seq)
156 + c.ack = p.header.seq
157 + c.diff = currentMicrosecond() - p.header.t
158 + c.state = state_connected
159 + c.closech = l.connch
160 + c.recvbuf = newPacketBuffer(window_size, int(p.header.seq))
161 + c.sendbuf = newPacketBuffer(window_size, seq)
162 +
163 + go c.loop()
164 + select {
165 + case c.recvch <- p:
166 + case <-c.recvchch:
167 + }
168 +
169 + l.conns[sid] = c
170 + ulog.Printf(2, "Listener(%v): New incoming connection #%d from %v (alive: %d)", l.conn.LocalAddr(), sid, addr, len(l.conns))
171 +
172 + l.accept <- c
173 + }
174 + }
175 +}
176 +
177 +func (l *UTPListener) Accept() (net.Conn, error) {
178 + return l.AcceptUTP()
179 +}
180 +
181 +func (l *UTPListener) AcceptUTP() (*UTPConn, error) {
182 + if l == nil || l.conn == nil {
183 + return nil, syscall.EINVAL
184 + }
185 + if l.lasterr != nil {
186 + return nil, l.lasterr
187 + }
188 + var timeout <-chan time.Time
189 + if !l.deadline.IsZero() {
190 + timeout = time.After(l.deadline.Sub(time.Now()))
191 + }
192 + select {
193 + case conn := <-l.accept:
194 + if conn == nil {
195 + return nil, errors.New("use of closed network connection")
196 + }
197 + return conn, nil
198 + case err := <-l.err:
199 + l.lasterr = err
200 + return nil, err
201 + case <-timeout:
202 + return nil, &timeoutError{}
203 + }
204 +}
205 +
206 +func (l *UTPListener) Addr() net.Addr {
207 + return &UTPAddr{Addr: l.conn.LocalAddr()}
208 +}
209 +
210 +func (l *UTPListener) Close() error {
211 + if l == nil || l.conn == nil {
212 + return syscall.EINVAL
213 + }
214 + l.closech <- 0
215 + l.RawConn.Close()
216 + return nil
217 +}
218 +
219 +func (l *UTPListener) SetDeadline(t time.Time) error {
220 + if l == nil || l.conn == nil {
221 + return syscall.EINVAL
222 + }
223 + l.deadline = t
224 + return nil
225 +}
226 +
227 +type rawIncoming struct {
228 + b []byte
229 + addr net.Addr
230 +}
231 +
232 +type rawConn struct {
233 + conn net.PacketConn
234 + rdeadline, wdeadline time.Time
235 + in chan rawIncoming
236 + closed chan int
237 +}
238 +
239 +func newRawConn(conn net.PacketConn) *rawConn {
240 + return &rawConn{
241 + conn: conn,
242 + in: make(chan rawIncoming, 100),
243 + closed: make(chan int),
244 + }
245 +}
246 +
247 +func (c *rawConn) ok() bool { return c != nil && c.conn != nil }
248 +
249 +func (c *rawConn) ReadFrom(b []byte) (n int, addr net.Addr, err error) {
250 + if !c.ok() {
251 + return 0, nil, syscall.EINVAL
252 + }
253 + select {
254 + case <-c.closed:
255 + return 0, nil, errors.New("use of closed network connection")
256 + default:
257 + }
258 + var timeout <-chan time.Time
259 + if !c.rdeadline.IsZero() {
260 + timeout = time.After(c.rdeadline.Sub(time.Now()))
261 + }
262 + select {
263 + case r := <-c.in:
264 + return copy(b, r.b), r.addr, nil
265 + case <-timeout:
266 + return 0, nil, &timeoutError{}
267 + }
268 +}
269 +
270 +func (c *rawConn) WriteTo(b []byte, addr net.Addr) (n int, err error) {
271 + if !c.ok() {
272 + return 0, syscall.EINVAL
273 + }
274 + select {
275 + case <-c.closed:
276 + return 0, errors.New("use of closed network connection")
277 + default:
278 + }
279 + return c.conn.WriteTo(b, addr)
280 +}
281 +
282 +func (c *rawConn) Close() error {
283 + if !c.ok() {
284 + return syscall.EINVAL
285 + }
286 + select {
287 + case <-c.closed:
288 + return errors.New("use of closed network connection")
289 + default:
290 + close(c.closed)
291 + }
292 + return nil
293 +}
294 +
295 +func (c *rawConn) LocalAddr() net.Addr {
296 + if !c.ok() {
297 + return nil
298 + }
299 + return c.conn.LocalAddr()
300 +}
301 +
302 +func (c *rawConn) SetDeadline(t time.Time) error {
303 + if !c.ok() {
304 + return syscall.EINVAL
305 + }
306 + if err := c.SetReadDeadline(t); err != nil {
307 + return err
308 + }
309 + if err := c.SetWriteDeadline(t); err != nil {
310 + return err
311 + }
312 + return nil
313 +}
314 +
315 +func (c *rawConn) SetReadDeadline(t time.Time) error {
316 + if !c.ok() {
317 + return syscall.EINVAL
318 + }
319 + c.rdeadline = t
320 + return nil
321 +}
322 +
323 +func (c *rawConn) SetWriteDeadline(t time.Time) error {
324 + if !c.ok() {
325 + return syscall.EINVAL
326 + }
327 + c.wdeadline = t
328 + return nil
329 +}
Godeps/_workspace/src/github.com/h2so5/utp/log.go new
+50
@@ -0,0 +1,50 @@
1 +package utp
2 +
3 +import (
4 + "log"
5 + "os"
6 + "strconv"
7 +)
8 +
9 +type logger struct {
10 + level int
11 +}
12 +
13 +var ulog *logger
14 +
15 +func init() {
16 + logenv := os.Getenv("GO_UTP_LOGGING")
17 +
18 + var level int
19 + if len(logenv) > 0 {
20 + l, err := strconv.Atoi(logenv)
21 + if err != nil {
22 + log.Print("warning: GO_UTP_LOGGING must be numeric")
23 + } else {
24 + level = l
25 + }
26 + }
27 +
28 + ulog = &logger{level}
29 +}
30 +
31 +func (l *logger) Print(level int, v ...interface{}) {
32 + if l.level < level {
33 + return
34 + }
35 + log.Print(v...)
36 +}
37 +
38 +func (l *logger) Printf(level int, format string, v ...interface{}) {
39 + if l.level < level {
40 + return
41 + }
42 + log.Printf(format, v...)
43 +}
44 +
45 +func (l *logger) Println(level int, v ...interface{}) {
46 + if l.level < level {
47 + return
48 + }
49 + log.Println(v...)
50 +}
Godeps/_workspace/src/github.com/h2so5/utp/packet.go new
+240
@@ -0,0 +1,240 @@
1 +package utp
2 +
3 +import (
4 + "bytes"
5 + "encoding/binary"
6 + "fmt"
7 + "io"
8 + "sync"
9 +)
10 +
11 +type header struct {
12 + typ, ver int
13 + id uint16
14 + t, diff, wnd uint32
15 + seq, ack uint16
16 +}
17 +
18 +type extension struct {
19 + typ int
20 + payload []byte
21 +}
22 +
23 +type packet struct {
24 + header header
25 + ext []extension
26 + payload []byte
27 +}
28 +
29 +type outgoingPacket struct {
30 + typ int
31 + ext []extension
32 + payload []byte
33 +}
34 +
35 +func (p *packet) MarshalBinary() ([]byte, error) {
36 + firstExt := ext_none
37 + if len(p.ext) > 0 {
38 + firstExt = p.ext[0].typ
39 + }
40 + buf := new(bytes.Buffer)
41 + var beforeExt = []interface{}{
42 + // | type | ver |
43 + uint8(((byte(p.header.typ) << 4) & 0xF0) | (byte(p.header.ver) & 0xF)),
44 + // | extension |
45 + uint8(firstExt),
46 + }
47 + var afterExt = []interface{}{
48 + // | connection_id |
49 + uint16(p.header.id),
50 + // | timestamp_microseconds |
51 + uint32(p.header.t),
52 + // | timestamp_difference_microseconds |
53 + uint32(p.header.diff),
54 + // | wnd_size |
55 + uint32(p.header.wnd),
56 + // | seq_nr |
57 + uint16(p.header.seq),
58 + // | ack_nr |
59 + uint16(p.header.ack),
60 + }
61 +
62 + for _, v := range beforeExt {
63 + err := binary.Write(buf, binary.BigEndian, v)
64 + if err != nil {
65 + return nil, err
66 + }
67 + }
68 +
69 + if len(p.ext) > 0 {
70 + for i, e := range p.ext {
71 + next := ext_none
72 + if i < len(p.ext)-1 {
73 + next = p.ext[i+1].typ
74 + }
75 + var ext = []interface{}{
76 + // | extension |
77 + uint8(next),
78 + // | len |
79 + uint8(len(e.payload)),
80 + }
81 + for _, v := range ext {
82 + err := binary.Write(buf, binary.BigEndian, v)
83 + if err != nil {
84 + return nil, err
85 + }
86 + }
87 + _, err := buf.Write(e.payload)
88 + if err != nil {
89 + return nil, err
90 + }
91 + }
92 + }
93 +
94 + for _, v := range afterExt {
95 + err := binary.Write(buf, binary.BigEndian, v)
96 + if err != nil {
97 + return nil, err
98 + }
99 + }
100 +
101 + _, err := buf.Write(p.payload)
102 + if err != nil {
103 + return nil, err
104 + }
105 + return buf.Bytes(), nil
106 +}
107 +
108 +func (p *packet) UnmarshalBinary(data []byte) error {
109 + p.ext = nil
110 + buf := bytes.NewReader(data)
111 + var tv, e uint8
112 +
113 + var beforeExt = []interface{}{
114 + // | type | ver |
115 + (*uint8)(&tv),
116 + // | extension |
117 + (*uint8)(&e),
118 + }
119 + for _, v := range beforeExt {
120 + err := binary.Read(buf, binary.BigEndian, v)
121 + if err != nil {
122 + return err
123 + }
124 + }
125 +
126 + for e != ext_none {
127 + currentExt := int(e)
128 + var l uint8
129 + var ext = []interface{}{
130 + // | extension |
131 + (*uint8)(&e),
132 + // | len |
133 + (*uint8)(&l),
134 + }
135 + for _, v := range ext {
136 + err := binary.Read(buf, binary.BigEndian, v)
137 + if err != nil {
138 + return err
139 + }
140 + }
141 + payload := make([]byte, l)
142 + size, err := buf.Read(payload[:])
143 + if err != nil {
144 + return err
145 + }
146 + if size != len(payload) {
147 + return io.EOF
148 + }
149 + p.ext = append(p.ext, extension{typ: currentExt, payload: payload})
150 + }
151 +
152 + var afterExt = []interface{}{
153 + // | connection_id |
154 + (*uint16)(&p.header.id),
155 + // | timestamp_microseconds |
156 + (*uint32)(&p.header.t),
157 + // | timestamp_difference_microseconds |
158 + (*uint32)(&p.header.diff),
159 + // | wnd_size |
160 + (*uint32)(&p.header.wnd),
161 + // | seq_nr |
162 + (*uint16)(&p.header.seq),
163 + // | ack_nr |
164 + (*uint16)(&p.header.ack),
165 + }
166 + for _, v := range afterExt {
167 + err := binary.Read(buf, binary.BigEndian, v)
168 + if err != nil {
169 + return err
170 + }
171 + }
172 +
173 + p.header.typ = int((tv >> 4) & 0xF)
174 + p.header.ver = int(tv & 0xF)
175 +
176 + l := buf.Len()
177 + if l > 0 {
178 + p.payload = p.payload[:l]
179 + _, err := buf.Read(p.payload[:])
180 + if err != nil {
181 + return err
182 + }
183 + }
184 +
185 + return nil
186 +}
187 +
188 +func (p packet) String() string {
189 + var s string = fmt.Sprintf("[%d ", p.header.id)
190 + switch p.header.typ {
191 + case st_data:
192 + s += "ST_DATA"
193 + case st_fin:
194 + s += "ST_FIN"
195 + case st_state:
196 + s += "ST_STATE"
197 + case st_reset:
198 + s += "ST_RESET"
199 + case st_syn:
200 + s += "ST_SYN"
201 + }
202 + s += fmt.Sprintf(" seq:%d ack:%d len:%d", p.header.seq, p.header.ack, len(p.payload))
203 + s += "]"
204 + return s
205 +}
206 +
207 +var globalPool packetPool
208 +
209 +type packetPool struct {
210 + root *packetPoolNode
211 + mutex sync.Mutex
212 +}
213 +
214 +type packetPoolNode struct {
215 + p *packet
216 + next *packetPoolNode
217 +}
218 +
219 +func (o *packetPool) get() *packet {
220 + o.mutex.Lock()
221 + defer o.mutex.Unlock()
222 + r := o.root
223 + if r != nil {
224 + o.root = o.root.next
225 + return r.p
226 + } else {
227 + return &packet{
228 + payload: make([]byte, 0, mss),
229 + }
230 + }
231 +}
232 +
233 +func (o *packetPool) put(p *packet) {
234 + o.mutex.Lock()
235 + defer o.mutex.Unlock()
236 + o.root = &packetPoolNode{
237 + p: p,
238 + next: o.root,
239 + }
240 +}
Godeps/_workspace/src/github.com/h2so5/utp/ucat/.gitignore new
+3
@@ -0,0 +1,3 @@
1 +ucat
2 +random
3 +.trash/
Godeps/_workspace/src/github.com/h2so5/utp/ucat/Makefile new
+34
@@ -0,0 +1,34 @@
1 +# Run tests
2 +
3 +testnames=simple
4 +tests=$(addprefix test_, $(testnames))
5 +trash=.trash/
6 +
7 +all: ucat
8 +
9 +test: clean ucat ${tests}
10 + @echo ${tests}
11 + @echo "*** tests passed ***"
12 +
13 +# not sue why this doesn't work:
14 +# test_%: test_%.sh
15 +test_simple: test_simple.sh
16 + mkdir -p ${trash}
17 + @echo "*** running $@ ***"
18 + ./$@.sh
19 +
20 +clean:
21 + @echo "*** $@ ***"
22 + -rm -r ${trash}
23 +
24 +deps: random ucat
25 +
26 +ucat:
27 + go build
28 +
29 +random:
30 + @echo "*** installing $@ ***"
31 + go get github.com/jbenet/go-random/random
32 + go build -o random github.com/jbenet/go-random/random
33 +
34 +.PHONY: clean ucat ${tests}
Godeps/_workspace/src/github.com/h2so5/utp/ucat/test_simple.sh new
+49
@@ -0,0 +1,49 @@
1 +#!/bin/sh
2 +
3 +set -e # exit on error
4 +# set -v # verbose
5 +
6 +log() {
7 + echo "--> $1"
8 +}
9 +
10 +test_send() {
11 + file=$1_
12 + count=$2
13 + addr=localhost:8765
14 +
15 + # generate random data
16 + log "generating $count bytes of random data"
17 + ./random $count $RANDOM > ${file}expected
18 +
19 + # dialer sends
20 + log "sending from dialer"
21 + ./ucat -v $addr 2>&1 <${file}expected | sed "s/^/ dialer1: /" &
22 + ./ucat -v -l $addr 2>&1 >${file}actual1 | sed "s/^/listener1: /"
23 + diff ${file}expected ${file}actual1
24 + if test $? != 0; then
25 + log "sending from dialer failed. compare with:\n"
26 + log "diff ${file}expected ${file}actual1"
27 + exit 1
28 + fi
29 +
30 + # listener sends
31 + log "sending from listener"
32 + ./ucat -v -l $addr 2>&1 <${file}expected | sed "s/^/listener2: /" &
33 + ./ucat -v $addr 2>&1 >${file}actual2 | sed "s/^/ dialer2: /"
34 + diff ${file}expected ${file}actual2
35 + if test $? != 0; then
36 + log "sending from listener failed. compare with:\n"
37 + log "diff ${file}expected ${file}actual2"
38 + exit 1
39 + fi
40 +
41 + echo rm ${file}{expected,actual1,actual2}
42 + rm ${file}{expected,actual1,actual2}
43 + return 0
44 +}
45 +
46 +
47 +test_send ".trash/1KB" 1024
48 +test_send ".trash/1MB" 1048576
49 +test_send ".trash/1GB" 1073741824
Godeps/_workspace/src/github.com/h2so5/utp/ucat/ucat.go new
+188
@@ -0,0 +1,188 @@
1 +// package ucat provides an implementation of netcat using the go utp package.
2 +// It is meant to exercise the utp implementation.
3 +// Usage:
4 +// ucat [<local address>] <remote address>
5 +// ucat -l <local address>
6 +//
7 +// Address format is: [host]:port
8 +//
9 +// Note that uTP's congestion control gives priority to tcp flows (web traffic),
10 +// so you could use this ucat tool to transfer massive files without hogging
11 +// all the bandwidth.
12 +package main
13 +
14 +import (
15 + "flag"
16 + "fmt"
17 + "io"
18 + "net"
19 + "os"
20 + "os/signal"
21 + "syscall"
22 +
23 + utp "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/h2so5/utp"
24 +)
25 +
26 +var verbose = false
27 +
28 +// Usage prints out the usage of this module.
29 +// Assumes flags use go stdlib flag pacakage.
30 +var Usage = func() {
31 + text := `ucat - uTP netcat in Go
32 +
33 +Usage:
34 +
35 + listen: %s [<local address>] <remote address>
36 + dial: %s -l <local address>
37 +
38 +Address format is Go's: [host]:port
39 +`
40 +
41 + fmt.Fprintf(os.Stderr, text, os.Args[0], os.Args[0])
42 + flag.PrintDefaults()
43 +}
44 +
45 +type args struct {
46 + listen bool
47 + verbose bool
48 + localAddr string
49 + remoteAddr string
50 +}
51 +
52 +func parseArgs() args {
53 + var a args
54 +
55 + // setup + parse flags
56 + flag.BoolVar(&a.listen, "listen", false, "listen for connections")
57 + flag.BoolVar(&a.listen, "l", false, "listen for connections (short)")
58 + flag.BoolVar(&a.verbose, "v", false, "verbose debugging")
59 + flag.Usage = Usage
60 + flag.Parse()
61 + osArgs := flag.Args()
62 +
63 + if len(osArgs) < 1 {
64 + exit("")
65 + }
66 +
67 + if a.listen {
68 + a.localAddr = osArgs[0]
69 + } else {
70 + if len(osArgs) > 1 {
71 + a.localAddr = osArgs[0]
72 + a.remoteAddr = osArgs[1]
73 + } else {
74 + a.remoteAddr = osArgs[0]
75 + }
76 + }
77 +
78 + return a
79 +}
80 +
81 +func main() {
82 + args := parseArgs()
83 + verbose = args.verbose
84 +
85 + var err error
86 + if args.listen {
87 + err = Listen(args.localAddr)
88 + } else {
89 + err = Dial(args.localAddr, args.remoteAddr)
90 + }
91 +
92 + if err != nil {
93 + exit("%s", err)
94 + }
95 +}
96 +
97 +func exit(format string, vals ...interface{}) {
98 + if format != "" {
99 + fmt.Fprintf(os.Stderr, "ucat error: "+format+"\n", vals...)
100 + }
101 + Usage()
102 + os.Exit(1)
103 +}
104 +
105 +func log(format string, vals ...interface{}) {
106 + if verbose {
107 + fmt.Fprintf(os.Stderr, "ucat log: "+format+"\n", vals...)
108 + }
109 +}
110 +
111 +// Listen listens and accepts one incoming uTP connection on a given port,
112 +// and pipes all incoming data to os.Stdout.
113 +func Listen(localAddr string) error {
114 + l, err := utp.Listen("utp", localAddr)
115 + if err != nil {
116 + return err
117 + }
118 + log("listening at %s", l.Addr())
119 +
120 + c, err := l.Accept()
121 + if err != nil {
122 + return err
123 + }
124 + log("accepted connection from %s", c.RemoteAddr())
125 +
126 + // should be able to close listener here, but utp.Listener.Close
127 + // closes all open connections.
128 + defer l.Close()
129 +
130 + netcat(c)
131 + return c.Close()
132 +}
133 +
134 +// Dial connects to a remote address and pipes all os.Stdin to the remote end.
135 +// If localAddr is set, uses it to Dial from.
136 +func Dial(localAddr, remoteAddr string) error {
137 +
138 + var laddr net.Addr
139 + var err error
140 + if localAddr != "" {
141 + laddr, err = utp.ResolveUTPAddr("utp", localAddr)
142 + if err != nil {
143 + return fmt.Errorf("failed to resolve address %s", localAddr)
144 + }
145 + }
146 +
147 + if laddr != nil {
148 + log("dialing %s from %s", remoteAddr, laddr)
149 + } else {
150 + log("dialing %s", remoteAddr)
151 + }
152 +
153 + d := utp.Dialer{LocalAddr: laddr}
154 + c, err := d.Dial("utp", remoteAddr)
155 + if err != nil {
156 + return err
157 + }
158 + log("connected to %s", c.RemoteAddr())
159 +
160 + netcat(c)
161 + return c.Close()
162 +}
163 +
164 +func netcat(c net.Conn) {
165 + log("piping stdio to connection")
166 +
167 + done := make(chan struct{})
168 +
169 + go func() {
170 + n, _ := io.Copy(c, os.Stdin)
171 + log("sent %d bytes", n)
172 + done <- struct{}{}
173 + }()
174 + go func() {
175 + n, _ := io.Copy(os.Stdout, c)
176 + log("received %d bytes", n)
177 + done <- struct{}{}
178 + }()
179 +
180 + // wait until we exit.
181 + sigc := make(chan os.Signal, 1)
182 + signal.Notify(sigc, syscall.SIGHUP, syscall.SIGINT,
183 + syscall.SIGTERM, syscall.SIGQUIT)
184 + select {
185 + case <-done:
186 + case <-sigc:
187 + }
188 +}
Godeps/_workspace/src/github.com/h2so5/utp/utp.go new
+29
@@ -0,0 +1,29 @@
1 +package utp
2 +
3 +import "time"
4 +
5 +const (
6 + version = 1
7 +
8 + st_data = 0
9 + st_fin = 1
10 + st_state = 2
11 + st_reset = 3
12 + st_syn = 4
13 +
14 + ext_none = 0
15 + ext_selective_ack = 1
16 +
17 + header_size = 20
18 + mtu = 3200
19 + mss = mtu - header_size
20 + window_size = 100
21 +
22 + reset_timeout = time.Second
23 +)
24 +
25 +type timeoutError struct{}
26 +
27 +func (e *timeoutError) Error() string { return "i/o timeout" }
28 +func (e *timeoutError) Timeout() bool { return true }
29 +func (e *timeoutError) Temporary() bool { return true }
Godeps/_workspace/src/github.com/h2so5/utp/utp_test.go new
+601
@@ -0,0 +1,601 @@
1 +package utp
2 +
3 +import (
4 + "bytes"
5 + "io"
6 + "io/ioutil"
7 + "math"
8 + "math/rand"
9 + "net"
10 + "reflect"
11 + "testing"
12 + "time"
13 +)
14 +
15 +func init() {
16 + rand.Seed(time.Now().Unix())
17 +}
18 +
19 +func TestReadWrite(t *testing.T) {
20 + ln, err := Listen("utp", "127.0.0.1:0")
21 + if err != nil {
22 + t.Fatal(err)
23 + }
24 +
25 + raddr, err := ResolveUTPAddr("utp", ln.Addr().String())
26 + if err != nil {
27 + t.Fatal(err)
28 + }
29 +
30 + c, err := DialUTPTimeout("utp", nil, raddr, 1000*time.Millisecond)
31 + if err != nil {
32 + t.Fatal(err)
33 + }
34 + defer c.Close()
35 +
36 + err = ln.SetDeadline(time.Now().Add(1000 * time.Millisecond))
37 + if err != nil {
38 + t.Fatal(err)
39 + }
40 +
41 + s, err := ln.Accept()
42 + if err != nil {
43 + t.Fatal(err)
44 + }
45 + ln.Close()
46 +
47 + payload := []byte("Hello!")
48 + _, err = c.Write(payload)
49 + if err != nil {
50 + t.Fatal(err)
51 + }
52 +
53 + err = s.SetDeadline(time.Now().Add(1000 * time.Millisecond))
54 + if err != nil {
55 + t.Fatal(err)
56 + }
57 +
58 + var buf [256]byte
59 + l, err := s.Read(buf[:])
60 + if err != nil {
61 + t.Fatal(err)
62 + }
63 +
64 + if !bytes.Equal(payload, buf[:l]) {
65 + t.Errorf("expected payload of %v; got %v", payload, buf[:l])
66 + }
67 +
68 + payload2 := []byte("World!")
69 + _, err = s.Write(payload2)
70 + if err != nil {
71 + t.Fatal(err)
72 + }
73 +
74 + err = c.SetDeadline(time.Now().Add(1000 * time.Millisecond))
75 + if err != nil {
76 + t.Fatal(err)
77 + }
78 +
79 + l, err = c.Read(buf[:])
80 + if err != nil {
81 + t.Fatal(err)
82 + }
83 +
84 + if !bytes.Equal(payload2, buf[:l]) {
85 + t.Errorf("expected payload of %v; got %v", payload2, buf[:l])
86 + }
87 +}
88 +
89 +func TestRawReadWrite(t *testing.T) {
90 + ln, err := Listen("utp", "127.0.0.1:0")
91 + if err != nil {
92 + t.Fatal(err)
93 + }
94 + defer ln.Close()
95 +
96 + raddr, err := net.ResolveUDPAddr("udp", ln.Addr().String())
97 + if err != nil {
98 + t.Fatal(err)
99 + }
100 +
101 + c, err := net.DialUDP("udp", nil, raddr)
102 + if err != nil {
103 + t.Fatal(err)
104 + }
105 + defer c.Close()
106 +
107 + payload := []byte("Hello!")
108 + _, err = c.Write(payload)
109 + if err != nil {
110 + t.Fatal(err)
111 + }
112 +
113 + var buf [256]byte
114 + n, addr, err := ln.RawConn.ReadFrom(buf[:])
115 + if !bytes.Equal(payload, buf[:n]) {
116 + t.Errorf("expected payload of %v; got %v", payload, buf[:n])
117 + }
118 + if addr.String() != c.LocalAddr().String() {
119 + t.Errorf("expected addr of %v; got %v", c.LocalAddr(), addr.String())
120 + }
121 +}
122 +
123 +func TestLongReadWriteC2S(t *testing.T) {
124 + ln, err := Listen("utp", "127.0.0.1:0")
125 + if err != nil {
126 + t.Fatal(err)
127 + }
128 +
129 + raddr, err := ResolveUTPAddr("utp", ln.Addr().String())
130 + if err != nil {
131 + t.Fatal(err)
132 + }
133 +
134 + c, err := DialUTPTimeout("utp", nil, raddr, 1000*time.Millisecond)
135 + if err != nil {
136 + t.Fatal(err)
137 + }
138 + defer c.Close()
139 +
140 + err = ln.SetDeadline(time.Now().Add(1000 * time.Millisecond))
141 + if err != nil {
142 + t.Fatal(err)
143 + }
144 +
145 + s, err := ln.Accept()
146 + if err != nil {
147 + t.Fatal(err)
148 + }
149 + defer s.Close()
150 + ln.Close()
151 +
152 + var payload [10485760]byte
153 + for i := range payload {
154 + payload[i] = byte(rand.Int())
155 + }
156 +
157 + rch := make(chan []byte)
158 + ech := make(chan error, 2)
159 +
160 + go func() {
161 + defer c.Close()
162 + _, err := c.Write(payload[:])
163 + if err != nil {
164 + ech <- err
165 + }
166 + }()
167 +
168 + go func() {
169 + b, err := ioutil.ReadAll(s)
170 + if err != nil {
171 + ech <- err
172 + rch <- nil
173 + } else {
174 + ech <- nil
175 + rch <- b
176 + }
177 + }()
178 +
179 + err = <-ech
180 + if err != nil {
181 + t.Fatal(err)
182 + }
183 +
184 + r := <-rch
185 + if r == nil {
186 + return
187 + }
188 +
189 + if !bytes.Equal(r, payload[:]) {
190 + t.Errorf("expected payload of %d; got %d", len(payload[:]), len(r))
191 + }
192 +}
193 +
194 +func TestLongReadWriteS2C(t *testing.T) {
195 + ln, err := Listen("utp", "127.0.0.1:0")
196 + if err != nil {
197 + t.Fatal(err)
198 + }
199 +
200 + raddr, err := ResolveUTPAddr("utp", ln.Addr().String())
201 + if err != nil {
202 + t.Fatal(err)
203 + }
204 +
205 + c, err := DialUTPTimeout("utp", nil, raddr, 1000*time.Millisecond)
206 + if err != nil {
207 + t.Fatal(err)
208 + }
209 + defer c.Close()
210 +
211 + err = ln.SetDeadline(time.Now().Add(1000 * time.Millisecond))
212 + if err != nil {
213 + t.Fatal(err)
214 + }
215 +
216 + s, err := ln.Accept()
217 + if err != nil {
218 + t.Fatal(err)
219 + }
220 + defer s.Close()
221 + ln.Close()
222 +
223 + var payload [10485760]byte
224 + for i := range payload {
225 + payload[i] = byte(rand.Int())
226 + }
227 +
228 + rch := make(chan []byte)
229 + ech := make(chan error, 2)
230 +
231 + go func() {
232 + defer s.Close()
233 + _, err := s.Write(payload[:])
234 + if err != nil {
235 + ech <- err
236 + }
237 + }()
238 +
239 + go func() {
240 + b, err := ioutil.ReadAll(c)
241 + if err != nil {
242 + ech <- err
243 + rch <- nil
244 + } else {
245 + ech <- nil
246 + rch <- b
247 + }
248 + }()
249 +
250 + err = <-ech
251 + if err != nil {
252 + t.Fatal(err)
253 + }
254 +
255 + r := <-rch
256 + if r == nil {
257 + return
258 + }
259 +
260 + if !bytes.Equal(r, payload[:]) {
261 + t.Errorf("expected payload of %d; got %d", len(payload[:]), len(r))
262 + }
263 +}
264 +
265 +func TestAccept(t *testing.T) {
266 + ln, err := Listen("utp", "127.0.0.1:0")
267 + if err != nil {
268 + t.Fatal(err)
269 + }
270 + defer ln.Close()
271 +
272 + c, err := DialUTPTimeout("utp", nil, ln.Addr().(*UTPAddr), 200*time.Millisecond)
273 + if err != nil {
274 + t.Fatal(err)
275 + }
276 + defer c.Close()
277 +
278 + err = ln.SetDeadline(time.Now().Add(100 * time.Millisecond))
279 + _, err = ln.Accept()
280 + if err != nil {
281 + t.Fatal(err)
282 + }
283 +}
284 +
285 +func TestAcceptDeadline(t *testing.T) {
286 + ln, err := Listen("utp", "127.0.0.1:0")
287 + if err != nil {
288 + t.Fatal(err)
289 + }
290 + defer ln.Close()
291 + err = ln.SetDeadline(time.Now().Add(time.Millisecond))
292 + _, err = ln.Accept()
293 + if err == nil {
294 + t.Fatal("Accept should failed")
295 + }
296 +}
297 +
298 +func TestAcceptClosedListener(t *testing.T) {
299 + ln, err := Listen("utp", "127.0.0.1:0")
300 + if err != nil {
301 + t.Fatal(err)
302 + }
303 + err = ln.Close()
304 + if err != nil {
305 + t.Fatal(err)
306 + }
307 + _, err = ln.Accept()
308 + if err == nil {
309 + t.Fatal("Accept should failed")
310 + }
311 + _, err = ln.Accept()
312 + if err == nil {
313 + t.Fatal("Accept should failed")
314 + }
315 +}
316 +
317 +func TestDialer(t *testing.T) {
318 + ln, err := Listen("utp", "127.0.0.1:0")
319 + if err != nil {
320 + t.Fatal(err)
321 + }
322 + defer ln.Close()
323 +
324 + d := Dialer{}
325 + c, err := d.Dial("utp", ln.Addr().String())
326 + if err != nil {
327 + t.Fatal(err)
328 + }
329 + defer c.Close()
330 +}
331 +
332 +func TestDialerAddrs(t *testing.T) {
333 + ln, err := Listen("utp", "127.0.0.1:0")
334 + if err != nil {
335 + t.Fatal(err)
336 + }
337 + defer ln.Close()
338 +
339 + laddr, err := ResolveUTPAddr("utp", "127.0.0.1:45678")
340 + if err != nil {
341 + t.Fatal(err)
342 + }
343 +
344 + d := Dialer{LocalAddr: laddr}
345 + c1, err := d.Dial("utp", ln.Addr().String())
346 + if err != nil {
347 + t.Fatal(err)
348 + }
349 + defer c1.Close()
350 +
351 + c2, err := ln.Accept()
352 + if err != nil {
353 + t.Fatal(err)
354 + }
355 + defer c2.Close()
356 +
357 + eq := func(a, b net.Addr) bool {
358 + return a.String() == b.String()
359 + }
360 +
361 + if !eq(d.LocalAddr, c2.RemoteAddr()) {
362 + t.Fatal("dialer.LocalAddr not equal to c2.RemoteAddr ")
363 + }
364 + if !eq(c1.LocalAddr(), c2.RemoteAddr()) {
365 + t.Fatal("c1.LocalAddr not equal to c2.RemoteAddr ")
366 + }
367 + if !eq(c2.LocalAddr(), c1.RemoteAddr()) {
368 + t.Fatal("c2.LocalAddr not equal to c1.RemoteAddr ")
369 + }
370 +}
371 +
372 +func TestDialerTimeout(t *testing.T) {
373 + timeout := time.Millisecond * 200
374 + d := Dialer{Timeout: timeout}
375 + done := make(chan struct{})
376 +
377 + go func() {
378 + _, err := d.Dial("utp", "127.0.0.1:34567")
379 + if err == nil {
380 + t.Fatal("should not connect")
381 + }
382 + done <- struct{}{}
383 + }()
384 +
385 + select {
386 + case <-time.After(timeout * 2):
387 + t.Fatal("should have ended already")
388 + case <-done:
389 + }
390 +}
391 +
392 +func TestPacketBinary(t *testing.T) {
393 + h := header{
394 + typ: st_fin,
395 + ver: version,
396 + id: 100,
397 + t: 50000,
398 + diff: 10000,
399 + wnd: 65535,
400 + seq: 100,
401 + ack: 200,
402 + }
403 +
404 + e := []extension{
405 + extension{
406 + typ: ext_selective_ack,
407 + payload: []byte{0, 1, 0, 1},
408 + },
409 + extension{
410 + typ: ext_selective_ack,
411 + payload: []byte{100, 0, 200, 0},
412 + },
413 + }
414 +
415 + p := packet{
416 + header: h,
417 + ext: e,
418 + payload: []byte("abcdefg"),
419 + }
420 +
421 + b, err := p.MarshalBinary()
422 + if err != nil {
423 + t.Fatal(err)
424 + }
425 +
426 + p2 := packet{payload: make([]byte, 0, mss)}
427 + err = p2.UnmarshalBinary(b)
428 + if err != nil {
429 + t.Fatal(err)
430 + }
431 +
432 + if !reflect.DeepEqual(p, p2) {
433 + t.Errorf("expected packet of %v; got %v", p, p2)
434 + }
435 +}
436 +
437 +func TestUnmarshalShortPacket(t *testing.T) {
438 + b := make([]byte, 18)
439 + p := packet{}
440 + err := p.UnmarshalBinary(b)
441 +
442 + if err == nil {
443 + t.Fatal("UnmarshalBinary should fail")
444 + } else if err != io.EOF {
445 + t.Fatal(err)
446 + }
447 +}
448 +
449 +func TestWriteOnClosedChannel(t *testing.T) {
450 + ln, err := Listen("utp", "127.0.0.1:0")
451 + if err != nil {
452 + t.Fatal(err)
453 + }
454 + defer ln.Close()
455 +
456 + c, err := DialUTPTimeout("utp", nil, ln.Addr().(*UTPAddr), 200*time.Millisecond)
457 + if err != nil {
458 + t.Fatal(err)
459 + }
460 +
461 + go func() {
462 + for {
463 + _, err := c.Write([]byte{100})
464 + if err != nil {
465 + return
466 + }
467 + }
468 + }()
469 +
470 + c.Close()
471 +}
472 +
473 +func TestReadOnClosedChannel(t *testing.T) {
474 + ln, err := Listen("utp", "127.0.0.1:0")
475 + if err != nil {
476 + t.Fatal(err)
477 + }
478 + defer ln.Close()
479 +
480 + c, err := DialUTPTimeout("utp", nil, ln.Addr().(*UTPAddr), 200*time.Millisecond)
481 + if err != nil {
482 + t.Fatal(err)
483 + }
484 +
485 + go func() {
486 + for {
487 + var buf [16]byte
488 + _, err := c.Read(buf[:])
489 + if err != nil {
490 + return
491 + }
492 + }
493 + }()
494 +
495 + c.Close()
496 +}
497 +
498 +func TestPacketBuffer(t *testing.T) {
499 + size := 12
500 + b := newPacketBuffer(12, 1)
501 +
502 + if b.space() != size {
503 + t.Errorf("expected space == %d; got %d", size, b.space())
504 + }
505 +
506 + for i := 1; i <= size; i++ {
507 + b.push(&packet{header: header{seq: uint16(i)}})
508 + }
509 +
510 + if b.space() != 0 {
511 + t.Errorf("expected space == 0; got %d", b.space())
512 + }
513 +
514 + a := []byte{255, 7}
515 + ack := b.generateSelectiveACK()
516 + if !bytes.Equal(a, ack) {
517 + t.Errorf("expected ack == %v; got %v", a, ack)
518 + }
519 +
520 + err := b.push(&packet{header: header{seq: 15}})
521 + if err == nil {
522 + t.Fatal("push should fail")
523 + }
524 +
525 + all := b.all()
526 + if len(all) != size {
527 + t.Errorf("expected %d packets sequence; got %d", size, len(all))
528 + }
529 +
530 + f := b.fetch(6)
531 + if f == nil {
532 + t.Fatal("fetch should not fail")
533 + }
534 +
535 + b.compact()
536 +
537 + err = b.push(&packet{header: header{seq: 15}})
538 + if err != nil {
539 + t.Fatal(err)
540 + }
541 +
542 + err = b.push(&packet{header: header{seq: 17}})
543 + if err != nil {
544 + t.Fatal(err)
545 + }
546 +
547 + for i := 7; i <= size; i++ {
548 + f := b.fetch(uint16(i))
549 + if f == nil {
550 + t.Fatal("fetch should not fail")
551 + }
552 + }
553 +
554 + a = []byte{128, 2}
555 + ack = b.generateSelectiveACK()
556 + if !bytes.Equal(a, ack) {
557 + t.Errorf("expected ack == %v; got %v", a, ack)
558 + }
559 +
560 + all = b.all()
561 + if len(all) != 2 {
562 + t.Errorf("expected 2 packets sequence; got %d", len(all))
563 + }
564 +
565 + b.compact()
566 + if b.space() != 9 {
567 + t.Errorf("expected space == 9; got %d", b.space())
568 + }
569 +
570 + ack = b.generateSelectiveACK()
571 + b.processSelectiveACK(ack)
572 +
573 + all = b.all()
574 + if len(all) != 1 {
575 + t.Errorf("expected size == 1; got %d", len(all))
576 + }
577 +}
578 +
579 +func TestPacketBufferBoundary(t *testing.T) {
580 + begin := math.MaxUint16 - 3
581 + b := newPacketBuffer(12, begin)
582 + for i := begin; i != 5; i = (i + 1) % (math.MaxUint16 + 1) {
583 + err := b.push(&packet{header: header{seq: uint16(i)}})
584 + if err != nil {
585 + t.Fatal(err)
586 + }
587 + }
588 +}
589 +
590 +func TestTimedBufferNode(t *testing.T) {
591 + b := timedBuffer{d: time.Millisecond * 100}
592 + b.push(100)
593 + b.push(200)
594 + time.Sleep(time.Millisecond * 200)
595 + b.push(300)
596 + b.push(400)
597 + m := b.min()
598 + if m != 300 {
599 + t.Errorf("expected min == 300; got %d", m)
600 + }
601 +}
Godeps/_workspace/src/github.com/jbenet/go-multiaddr-net/convert.go
+2 -1
@@ -5,7 +5,7 @@ import (
5 "net"
6 "strings"
7
8 - utp "github.com/h2so5/utp"
8 + utp "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/h2so5/utp"
9 ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
10 )
11
@@ -154,6 +154,7 @@ func DialArgs(m ma.Multiaddr) (string, string, error) {
154 if parts[2] == "udp" && len(parts) > 4 && parts[4] == "utp" {
155 network = parts[4]
156 }
157 +
158 var host string
159 switch parts[0] {
160 case "ip4":
Godeps/_workspace/src/github.com/jbenet/go-multiaddr-net/convert_test.go
+1 -1
@@ -4,7 +4,7 @@ import (
4 "net"
5 "testing"
6
7 - utp "github.com/h2so5/utp"
7 + utp "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/h2so5/utp"
8 ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
9 )
10
Godeps/_workspace/src/github.com/jbenet/go-multiaddr-net/net.go
+5 -8
@@ -4,7 +4,7 @@ import (
4 "fmt"
5 "net"
6
7 - utp "github.com/h2so5/utp"
7 + utp "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/h2so5/utp"
8 ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
9 )
10
@@ -226,14 +226,11 @@ func Listen(laddr ma.Multiaddr) (Listener, error) {
226 switch lnet {
227 case "utp":
228 nl, err = utp.Listen(lnet, lnaddr)
229 - if err != nil {
230 - return nil, err
231 - }
232 - case "tcp":
229 + default:
230 nl, err = net.Listen(lnet, lnaddr)
234 - if err != nil {
235 - return nil, err
236 - }
231 + }
232 + if err != nil {
233 + return nil, err
234 }
235
236 return &maListener{
Godeps/_workspace/src/github.com/jbenet/go-multiaddr-net/net_test.go
+30
@@ -138,6 +138,36 @@ func TestListen(t *testing.T) {
138 wg.Wait()
139 }
140
141 +func TestListenAddrs(t *testing.T) {
142 +
143 + test := func(addr string, succeed bool) {
144 +
145 + maddr := newMultiaddr(t, addr)
146 + l, err := Listen(maddr)
147 + if !succeed {
148 + if err == nil {
149 + t.Fatal("succeeded in listening", addr)
150 + }
151 + return
152 + }
153 + if succeed && err != nil {
154 + t.Fatal("failed to listen", addr, err)
155 + }
156 + if l == nil {
157 + t.Fatal("failed to listen", addr, succeed, err)
158 + }
159 +
160 + if err = l.Close(); err != nil {
161 + t.Fatal("failed to close listener", addr, err)
162 + }
163 + }
164 +
165 + test("/ip4/127.0.0.1/tcp/4324", true)
166 + test("/ip4/127.0.0.1/udp/4325", false)
167 + test("/ip4/127.0.0.1/udp/4326/udt", false)
168 + test("/ip4/127.0.0.1/udp/4326/utp", true)
169 +}
170 +
171 func TestListenAndDial(t *testing.T) {
172
173 maddr := newMultiaddr(t, "/ip4/127.0.0.1/tcp/4323")