Compare commits

..
Author SHA1 Message Date
xtaci 13e980cdff lint 2016-06-30 14:45:26 +08:00
11 changed files with 301 additions and 69 deletions
+16
View File
@@ -0,0 +1,16 @@
language: go
go:
- 1.6
before_install:
- go get github.com/mattn/goveralls
- go get golang.org/x/tools/cmd/cover
install:
- go get github.com/xtaci/kcptun/client
- go get github.com/xtaci/kcptun/server
before_script:
script:
- cd $HOME/gopath/src/github.com/xtaci/kcptun/client
- $HOME/gopath/bin/goveralls -service=travis-ci
- cd $HOME/gopath/src/github.com/xtaci/kcptun/server
- $HOME/gopath/bin/goveralls -service=travis-ci
- exit 0
+85 -22
View File
@@ -1,9 +1,37 @@
# *kcptun*
TCP流转换为KCP+UDP流,:zap:***[下载地址](https://github.com/xtaci/kcptun/releases/latest)***:zap:工作示意图:
# <img src="logo.png" alt="kcptun" height="60px" />
[![GoDoc][1]][2] [![Release][13]][14] [![Powered][17]][18] [![Build Status][3]][4] [![Go Report Card][5]][6] [![Downloads][15]][16]
[1]: https://godoc.org/github.com/xtaci/kcptun?status.svg
[2]: https://godoc.org/github.com/xtaci/kcptun
[3]: https://travis-ci.org/xtaci/kcptun.svg?branch=master
[4]: https://travis-ci.org/xtaci/kcptun
[5]: https://goreportcard.com/badge/github.com/xtaci/kcptun
[6]: https://goreportcard.com/report/github.com/xtaci/kcptun
[7]: https://img.shields.io/badge/license-MIT-blue.svg
[8]: https://raw.githubusercontent.com/xtaci/kcptun/master/LICENSE.md
[9]: https://img.shields.io/github/stars/xtaci/kcptun.svg
[10]: https://github.com/xtaci/kcptun/stargazers
[11]: https://img.shields.io/github/forks/xtaci/kcptun.svg
[12]: https://github.com/xtaci/kcptun/network
[13]: https://img.shields.io/github/release/xtaci/kcptun.svg
[14]: https://github.com/xtaci/kcptun/releases/latest
[15]: https://img.shields.io/github/downloads/xtaci/kcptun/total.svg?maxAge=2592000
[16]: https://github.com/xtaci/kcptun/releases
[17]: https://img.shields.io/badge/KCP-Powered-blue.svg
[18]: https://github.com/skywind3000/kcp
[19]: https://img.shields.io/docker/pulls/xtaci/kcptun.svg?maxAge=2592000
[20]: https://hub.docker.com/r/xtaci/kcptun/
TCP流转换为KCP+UDP流,:zap:***[官方下载地址](https://github.com/xtaci/kcptun/releases/latest)***:zap:工作示意图:
![kcptun](kcptun.png)
***kcptun是[kcp](https://github.com/xtaci/kcp-go)协议的一个简单应用,可以用于任意tcp网络程序的传输承载,以提高网络流畅度,降低掉线情况。***
***kcptun是[kcp-go](https://github.com/xtaci/kcp-go)的一个测试应用,可以用于任意tcp网络程序的传输承载(尤其用于游戏数据传输测试),用于优化丢包环境下的网络流畅度。***
### *快速设定* :lollipop:
```
服务器: ./server_linux_amd64 -t "127.0.0.1:1080" -l ":554" -mode fast2 // 转发到本地1080端口
客户端: ./client_darwin_amd64 -r "服务器IP地址:554" -l ":1080" -mode fast2 // 监听本地1080端口
```
### *使用の方法* :lollipop:
![client](client.png)
@@ -14,6 +42,23 @@ TCP流转换为KCP+UDP流,:zap:***[下载地址](https://github.com/xtaci/kcpt
2. 跨运营商的流量传输
3. 其他高丢包通信链路的TCP承载
### *推荐参数* :lollipop:
```
适用大部分ADSL接入(非对称上下行)的参数(实验环境电信100M ADSL)
SERVER: -mtu 1400 -sndwnd 2048 -rcvwnd 2048 -mode fast2
CLIENT: -mtu 1400 -sndwnd 256 -rcvwnd 2048 -mode fast2 -dscp 46
```
*简易调优方法*
> 第一步:同时在两端逐步增大client rcvwnd和server sndwnd;
> 第二步:尝试下载,观察如果带宽利用率接近物理带宽则停止,否则跳转到第一步。
*巭孬嫑乱动*
### *DSCP* :lollipop:
DSCP差分服务代码点(Differentiated Services Code Point),IETF于1998年12月发布了Diff-ServDifferentiated Service)的QoS分类标准。它在每个数据包IP头部的服务类别TOS标识字节中,利用已使用的6比特和未使用的2比特,通过编码值来区分优先级。
常用DSCP值可以参考[Wikipedia DSCP](https://en.wikipedia.org/wiki/Differentiated_services#Commonly_used_DSCP_values),至于有没有用,完全取决于数据包经过的设备。
### *内置模式* :lollipop:
响应速度:
*fast3 >* ***[fast2]*** *> fast > normal > default*
@@ -22,31 +67,49 @@ TCP流转换为KCP+UDP流,:zap:***[下载地址](https://github.com/xtaci/kcpt
中间mode参数比较均衡,总之就是越快越浪费带宽,推荐模式 ***fast2***
更高级的 ***手动档*** 需要理解KCP协议,并通过 ***隐藏参数*** 调整,例如:
```
-mode manual -nodelay 1 -resend 4 -nc 1 -interval 20 -fec 4
-mode manual -nodelay 1 -resend 2 -nc 1 -interval 20
```
### *前向纠错* :lollipop:
前向纠错采用Reed Solomon纠删码, 它的基本原理如下: 给定n个数据块d1, d2,…, dn,n和一个正整数m, RS根据n个数据块生成m个校验块, c1, c2,…, cm。 对于任意的n和m, 从n个原始数据块和m 个校验块中任取n块就能解码出原始数据, 即RS最多容忍m个数据块或者校验块同时丢失。
![reed-solomon](rs.png)
通过```-datashard 10 -parityshard 3``` 可以调整Reed Solomon参数。
### *Snappy数据流压缩* :lollipop:
> Snappy is a compression/decompression library. It does not aim for maximum
> compression, or compatibility with any other compression library; instead,
> it aims for very high speeds and reasonable compression. For instance,
> compared to the fastest mode of zlib, Snappy is an order of magnitude faster
> for most inputs, but the resulting compressed files are anywhere from 20% to
> 100% bigger.
Reference: http://google.github.io/snappy/
### *SNMP* :lollipop:
```go
// Snmp defines network statistics indicator
type Snmp struct {
BytesSent uint64 // payload bytes sent
BytesReceived uint64
MaxConn uint64
ActiveOpens uint64
PassiveOpens uint64
CurrEstab uint64
InErrs uint64
InCsumErrors uint64 // checksum errors
InSegs uint64
OutSegs uint64
OutBytes uint64 // udp bytes sent
RetransSegs uint64
FastRetransSegs uint64
LostSegs uint64
RepeatSegs uint64
FECRecovered uint64
FECErrs uint64
FECSegs uint64 // fec segments received
BytesSent uint64 // payload bytes sent
BytesReceived uint64
MaxConn uint64
ActiveOpens uint64
PassiveOpens uint64
CurrEstab uint64
InErrs uint64
InCsumErrors uint64 // checksum errors
InSegs uint64
OutSegs uint64
OutBytes uint64 // udp bytes sent
RetransSegs uint64
FastRetransSegs uint64
EarlyRetransSegs uint64
LostSegs uint64
RepeatSegs uint64
FECRecovered uint64
FECErrs uint64
FECSegs uint64 // fec segments received
}
```
BIN
View File
Binary file not shown.

Before

Width:  |  Height:  |  Size: 54 KiB

After

Width:  |  Height:  |  Size: 66 KiB

+116 -34
View File
@@ -1,6 +1,7 @@
package main
import (
"crypto/sha1"
"io"
"log"
"math/rand"
@@ -8,14 +9,50 @@ import (
"os"
"time"
"golang.org/x/crypto/pbkdf2"
"github.com/golang/snappy"
"github.com/hashicorp/yamux"
"github.com/urfave/cli"
"github.com/xtaci/kcp-go"
)
var VERSION = "SELFBUILD"
var (
// VERSION is injected by buildflags
VERSION = "SELFBUILD"
// SALT is use for pbkdf2 key expansion
SALT = "kcp-go"
)
func handleClient(p1, p2 net.Conn) {
type compStream struct {
conn net.Conn
w *snappy.Writer
r *snappy.Reader
}
func (c *compStream) Read(p []byte) (n int, err error) {
return c.r.Read(p)
}
func (c *compStream) Write(p []byte) (n int, err error) {
n, err = c.w.Write(p)
err = c.w.Flush()
return n, err
}
func (c *compStream) Close() error {
return c.conn.Close()
}
func newCompStream(conn net.Conn) *compStream {
c := new(compStream)
c.conn = conn
c.w = snappy.NewBufferedWriter(conn)
c.r = snappy.NewReader(conn)
return c
}
func handleClient(p1, p2 io.ReadWriteCloser) {
log.Println("stream opened")
defer log.Println("stream closed")
defer p1.Close()
@@ -71,11 +108,21 @@ func main() {
Usage: "key for communcation, must be the same as kcptun server",
EnvVar: "KCPTUN_KEY",
},
cli.StringFlag{
Name: "crypt",
Value: "aes",
Usage: "methods for encryption: aes, tea, xor, none",
},
cli.StringFlag{
Name: "mode",
Value: "fast",
Usage: "mode for communication: fast3, fast2, fast, normal",
},
cli.IntFlag{
Name: "conn",
Value: 1,
Usage: "establish N physical connections as specified by 'conn' to server",
},
cli.IntFlag{
Name: "mtu",
Value: 1350,
@@ -91,10 +138,19 @@ func main() {
Value: 1024,
Usage: "set receive window size(num of packets)",
},
cli.BoolFlag{
Name: "nocomp",
Usage: "disable compression",
},
cli.IntFlag{
Name: "fec",
Value: 4,
Usage: "set FEC group size, must be the same as server",
Name: "datashard",
Value: 10,
Usage: "set reed-solomon erasure coding - datashard",
},
cli.IntFlag{
Name: "parityshard",
Value: 3,
Usage: "set reed-solomon erasure coding - parityshard",
},
cli.BoolFlag{
Name: "acknodelay",
@@ -132,13 +188,9 @@ func main() {
checkError(err)
listener, err := net.ListenTCP("tcp", addr)
checkError(err)
log.Println("listening on:", listener.Addr())
START_KCP:
pass := pbkdf2.Key([]byte(c.String("key")), []byte(SALT), 4096, 32, sha1.New)
// kcp server
kcpconn, err := kcp.DialWithOptions(c.Int("fec"), c.String("remoteaddr"), []byte(c.String("key")))
checkError(err)
nodelay, interval, resend, nc := c.Int("nodelay"), c.Int("interval"), c.Int("resend"), c.Int("nc")
switch c.String("mode") {
@@ -152,48 +204,78 @@ func main() {
nodelay, interval, resend, nc = 1, 10, 2, 1
}
log.Println("listening on:", listener.Addr())
log.Println("encryption:", c.String("crypt"))
log.Println("nodelay parameters:", nodelay, interval, resend, nc)
log.Println("remote address:", c.String("remoteaddr"))
log.Println("sndwnd:", c.Int("sndwnd"), "rcvwnd:", c.Int("rcvwnd"))
log.Println("compression:", !c.Bool("nocomp"))
log.Println("mtu:", c.Int("mtu"))
log.Println("fec:", c.Int("fec"))
log.Println("datashard:", c.Int("datashard"), "parityshard:", c.Int("parityshard"))
log.Println("acknodelay:", c.Bool("acknodelay"))
log.Println("dscp:", c.Int("dscp"))
log.Println("conn:", c.Int("conn"))
kcpconn.SetNoDelay(nodelay, interval, resend, nc)
kcpconn.SetWindowSize(c.Int("sndwnd"), c.Int("rcvwnd"))
kcpconn.SetMtu(c.Int("mtu"))
kcpconn.SetACKNoDelay(c.Bool("acknodelay"))
kcpconn.SetDSCP(c.Int("dscp"))
createConn := func() *yamux.Session {
var block kcp.BlockCrypt
switch c.String("crypt") {
case "tea":
block, _ = kcp.NewTEABlockCrypt(pass[:16])
case "xor":
block, _ = kcp.NewSimpleXORBlockCrypt(pass)
case "none":
block, _ = kcp.NewNoneBlockCrypt(pass)
default:
block, _ = kcp.NewAESBlockCrypt(pass)
}
kcpconn, err := kcp.DialWithOptions(c.String("remoteaddr"), block, c.Int("datashard"), c.Int("parityshard"))
checkError(err)
kcpconn.SetNoDelay(nodelay, interval, resend, nc)
kcpconn.SetWindowSize(c.Int("sndwnd"), c.Int("rcvwnd"))
kcpconn.SetMtu(c.Int("mtu"))
kcpconn.SetACKNoDelay(c.Bool("acknodelay"))
kcpconn.SetDSCP(c.Int("dscp"))
// stream multiplex
var mux *yamux.Session
config := &yamux.Config{
AcceptBacklog: 256,
EnableKeepAlive: true,
KeepAliveInterval: 30 * time.Second,
ConnectionWriteTimeout: 30 * time.Second,
MaxStreamWindowSize: 16777216,
LogOutput: os.Stderr,
// stream multiplex
config := &yamux.Config{
AcceptBacklog: 256,
EnableKeepAlive: true,
KeepAliveInterval: 30 * time.Second,
ConnectionWriteTimeout: 30 * time.Second,
MaxStreamWindowSize: 16777216,
LogOutput: os.Stderr,
}
var session *yamux.Session
if c.Bool("nocomp") {
session, err = yamux.Client(kcpconn, config)
} else {
session, err = yamux.Client(newCompStream(kcpconn), config)
}
checkError(err)
return session
}
session, err := yamux.Client(kcpconn, config)
checkError(err)
mux = session
numconn := uint16(c.Int("conn"))
var muxes []*yamux.Session
for i := uint16(0); i < numconn; i++ {
muxes = append(muxes, createConn())
}
rr := uint16(0)
for {
p1, err := listener.AcceptTCP()
if err != nil {
log.Println(err)
continue
}
checkError(err)
mux := muxes[rr%numconn]
p2, err := mux.Open()
if err != nil { // yamux failure
log.Println(err)
kcpconn.Close()
p1.Close()
goto START_KCP
mux.Close()
muxes[rr%numconn] = createConn()
continue
}
go handleClient(p1, p2)
rr++
}
}
myApp.Run(os.Args)
+2 -2
View File
@@ -12,10 +12,10 @@ import (
)
func init() {
go sig_handler()
go sigHandler()
}
func sig_handler() {
func sigHandler() {
ch := make(chan os.Signal, 1)
signal.Notify(ch, syscall.SIGUSR1)
BIN
View File
Binary file not shown.

Before

Width:  |  Height:  |  Size: 22 KiB

After

Width:  |  Height:  |  Size: 20 KiB

BIN
View File
Binary file not shown.

After

Width:  |  Height:  |  Size: 4.5 KiB

BIN
View File
Binary file not shown.

After

Width:  |  Height:  |  Size: 25 KiB

BIN
View File
Binary file not shown.

Before

Width:  |  Height:  |  Size: 54 KiB

After

Width:  |  Height:  |  Size: 63 KiB

+80 -9
View File
@@ -1,6 +1,7 @@
package main
import (
"crypto/sha1"
"io"
"log"
"math/rand"
@@ -8,15 +9,51 @@ import (
"os"
"time"
"golang.org/x/crypto/pbkdf2"
"github.com/golang/snappy"
"github.com/hashicorp/yamux"
"github.com/urfave/cli"
"github.com/xtaci/kcp-go"
)
var VERSION = "SELFBUILD"
var (
// VERSION is injected by buildflags
VERSION = "SELFBUILD"
// SALT is use for pbkdf2 key expansion
SALT = "kcp-go"
)
type compStream struct {
conn net.Conn
w *snappy.Writer
r *snappy.Reader
}
func (c *compStream) Read(p []byte) (n int, err error) {
return c.r.Read(p)
}
func (c *compStream) Write(p []byte) (n int, err error) {
n, err = c.w.Write(p)
err = c.w.Flush()
return n, err
}
func (c *compStream) Close() error {
return c.conn.Close()
}
func newCompStream(conn net.Conn) *compStream {
c := new(compStream)
c.conn = conn
c.w = snappy.NewBufferedWriter(conn)
c.r = snappy.NewReader(conn)
return c
}
// handle multiplex-ed connection
func handleMux(conn *kcp.UDPSession, target string) {
func handleMux(conn io.ReadWriteCloser, target string) {
// stream multiplex
var mux *yamux.Session
config := &yamux.Config{
@@ -50,7 +87,7 @@ func handleMux(conn *kcp.UDPSession, target string) {
}
}
func handleClient(p1, p2 net.Conn) {
func handleClient(p1, p2 io.ReadWriteCloser) {
log.Println("stream opened")
defer log.Println("stream closed")
defer p1.Close()
@@ -99,6 +136,11 @@ func main() {
Usage: "key for communcation, must be the same as kcptun client",
EnvVar: "KCPTUN_KEY",
},
cli.StringFlag{
Name: "crypt",
Value: "aes",
Usage: "methods for encryption: aes, tea, xor, none",
},
cli.StringFlag{
Name: "mode",
Value: "fast",
@@ -119,10 +161,19 @@ func main() {
Value: 1024,
Usage: "set receive window size(num of packets)",
},
cli.BoolFlag{
Name: "nocomp",
Usage: "disable compression",
},
cli.IntFlag{
Name: "fec",
Value: 4,
Usage: "set FEC group size, must be the same as client",
Name: "datashard",
Value: 10,
Usage: "set reed-solomon erasure coding - datashard",
},
cli.IntFlag{
Name: "parityshard",
Value: 3,
Usage: "set reed-solomon erasure coding - parityshard",
},
cli.BoolFlag{
Name: "acknodelay",
@@ -168,15 +219,30 @@ func main() {
nodelay, interval, resend, nc = 1, 10, 2, 1
}
lis, err := kcp.ListenWithOptions(c.Int("fec"), c.String("listen"), []byte(c.String("key")))
pass := pbkdf2.Key([]byte(c.String("key")), []byte(SALT), 4096, 32, sha1.New)
var block kcp.BlockCrypt
switch c.String("crypt") {
case "tea":
block, _ = kcp.NewTEABlockCrypt(pass[:16])
case "xor":
block, _ = kcp.NewSimpleXORBlockCrypt(pass)
case "none":
block, _ = kcp.NewNoneBlockCrypt(pass)
default:
block, _ = kcp.NewAESBlockCrypt(pass)
}
lis, err := kcp.ListenWithOptions(c.String("listen"), block, c.Int("datashard"), c.Int("parityshard"))
if err != nil {
log.Fatal(err)
}
log.Println("listening on ", lis.Addr())
log.Println("encryption:", c.String("crypt"))
log.Println("nodelay parameters:", nodelay, interval, resend, nc)
log.Println("sndwnd:", c.Int("sndwnd"), "rcvwnd:", c.Int("rcvwnd"))
log.Println("compression:", !c.Bool("nocomp"))
log.Println("mtu:", c.Int("mtu"))
log.Println("fec:", c.Int("fec"))
log.Println("datashard:", c.Int("datashard"), "parityshard:", c.Int("parityshard"))
log.Println("acknodelay:", c.Bool("acknodelay"))
log.Println("dscp:", c.Int("dscp"))
for {
@@ -187,7 +253,12 @@ func main() {
conn.SetWindowSize(c.Int("sndwnd"), c.Int("rcvwnd"))
conn.SetACKNoDelay(c.Bool("acknodelay"))
conn.SetDSCP(c.Int("dscp"))
go handleMux(conn, c.String("target"))
if c.Bool("nocomp") {
go handleMux(conn, c.String("target"))
} else {
go handleMux(newCompStream(conn), c.String("target"))
}
} else {
log.Println(err)
}
+2 -2
View File
@@ -12,10 +12,10 @@ import (
)
func init() {
go sig_handler()
go sigHandler()
}
func sig_handler() {
func sigHandler() {
ch := make(chan os.Signal, 1)
signal.Notify(ch, syscall.SIGUSR1)