mirror of
https://github.com/v2fly/v2ray-core.git
synced 2026-06-17 16:29:55 -04:00
move transport methods from net to io
This commit is contained in:
@@ -11,6 +11,13 @@ func Release(buffer *Buffer) {
|
||||
}
|
||||
}
|
||||
|
||||
func Len(buffer *Buffer) int {
|
||||
if buffer == nil {
|
||||
return 0
|
||||
}
|
||||
return buffer.Len()
|
||||
}
|
||||
|
||||
// Buffer is a recyclable allocation of a byte array. Buffer.Release() recycles
|
||||
// the buffer into an internal buffer pool, in order to recreate a buffer more
|
||||
// quickly.
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
package crypto
|
||||
|
||||
type Authenticator interface {
|
||||
AuthBytes() int
|
||||
AuthSize() int
|
||||
Authenticate(auth []byte, data []byte) []byte
|
||||
}
|
||||
|
||||
122
common/io/reader.go
Normal file
122
common/io/reader.go
Normal file
@@ -0,0 +1,122 @@
|
||||
package io // import "github.com/v2ray/v2ray-core/common/io"
|
||||
|
||||
import (
|
||||
"io"
|
||||
|
||||
"github.com/v2ray/v2ray-core/common/alloc"
|
||||
"github.com/v2ray/v2ray-core/common/crypto"
|
||||
"github.com/v2ray/v2ray-core/common/serial"
|
||||
"github.com/v2ray/v2ray-core/transport"
|
||||
)
|
||||
|
||||
// ReadFrom reads from a reader and put all content to a buffer.
|
||||
// If buffer is nil, ReadFrom creates a new normal buffer.
|
||||
func ReadFrom(reader io.Reader, buffer *alloc.Buffer) (*alloc.Buffer, error) {
|
||||
if buffer == nil {
|
||||
buffer = alloc.NewBuffer()
|
||||
}
|
||||
nBytes, err := reader.Read(buffer.Value)
|
||||
buffer.Slice(0, nBytes)
|
||||
return buffer, err
|
||||
}
|
||||
|
||||
type Reader interface {
|
||||
Read() (*alloc.Buffer, error)
|
||||
}
|
||||
|
||||
type AdaptiveReader struct {
|
||||
reader io.Reader
|
||||
allocate func() *alloc.Buffer
|
||||
isLarge bool
|
||||
}
|
||||
|
||||
func NewAdaptiveReader(reader io.Reader) *AdaptiveReader {
|
||||
return &AdaptiveReader{
|
||||
reader: reader,
|
||||
allocate: alloc.NewBuffer,
|
||||
isLarge: false,
|
||||
}
|
||||
}
|
||||
|
||||
func (this *AdaptiveReader) Read() (*alloc.Buffer, error) {
|
||||
buffer, err := ReadFrom(this.reader, this.allocate())
|
||||
|
||||
if buffer.IsFull() && !this.isLarge {
|
||||
this.allocate = alloc.NewLargeBuffer
|
||||
this.isLarge = true
|
||||
} else if !buffer.IsFull() {
|
||||
this.allocate = alloc.NewBuffer
|
||||
this.isLarge = false
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
alloc.Release(buffer)
|
||||
return nil, err
|
||||
}
|
||||
return buffer, nil
|
||||
}
|
||||
|
||||
type ChunkReader struct {
|
||||
reader io.Reader
|
||||
}
|
||||
|
||||
func NewChunkReader(reader io.Reader) *ChunkReader {
|
||||
return &ChunkReader{
|
||||
reader: reader,
|
||||
}
|
||||
}
|
||||
|
||||
func (this *ChunkReader) Read() (*alloc.Buffer, error) {
|
||||
buffer := alloc.NewLargeBuffer()
|
||||
if _, err := io.ReadFull(this.reader, buffer.Value[:2]); err != nil {
|
||||
alloc.Release(buffer)
|
||||
return nil, err
|
||||
}
|
||||
length := serial.BytesLiteral(buffer.Value[:2]).Uint16Value()
|
||||
if _, err := io.ReadFull(this.reader, buffer.Value[:length]); err != nil {
|
||||
alloc.Release(buffer)
|
||||
return nil, err
|
||||
}
|
||||
buffer.Slice(0, int(length))
|
||||
return buffer, nil
|
||||
}
|
||||
|
||||
type AuthenticationReader struct {
|
||||
reader Reader
|
||||
authenticator crypto.Authenticator
|
||||
authBeforePayload bool
|
||||
}
|
||||
|
||||
func NewAuthenticationReader(reader io.Reader, auth crypto.Authenticator, authBeforePayload bool) *AuthenticationReader {
|
||||
return &AuthenticationReader{
|
||||
reader: NewChunkReader(reader),
|
||||
authenticator: auth,
|
||||
authBeforePayload: authBeforePayload,
|
||||
}
|
||||
}
|
||||
|
||||
func (this *AuthenticationReader) Read() (*alloc.Buffer, error) {
|
||||
buffer, err := this.reader.Read()
|
||||
if err != nil {
|
||||
alloc.Release(buffer)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
authSize := this.authenticator.AuthSize()
|
||||
var authBytes, payloadBytes []byte
|
||||
if this.authBeforePayload {
|
||||
authBytes = buffer.Value[:authSize]
|
||||
payloadBytes = buffer.Value[authSize:]
|
||||
} else {
|
||||
payloadBytes = buffer.Value[:authSize]
|
||||
authBytes = buffer.Value[authSize:]
|
||||
}
|
||||
|
||||
actualAuthBytes := this.authenticator.Authenticate(nil, payloadBytes)
|
||||
if !serial.BytesLiteral(authBytes).Equals(serial.BytesLiteral(actualAuthBytes)) {
|
||||
alloc.Release(buffer)
|
||||
return nil, transport.CorruptedPacket
|
||||
}
|
||||
buffer.Value = payloadBytes
|
||||
return buffer, nil
|
||||
}
|
||||
42
common/io/transport.go
Normal file
42
common/io/transport.go
Normal file
@@ -0,0 +1,42 @@
|
||||
package io
|
||||
|
||||
import (
|
||||
"io"
|
||||
|
||||
"github.com/v2ray/v2ray-core/common/alloc"
|
||||
)
|
||||
|
||||
func RawReaderToChan(stream chan<- *alloc.Buffer, reader io.Reader) error {
|
||||
return ReaderToChan(stream, NewAdaptiveReader(reader))
|
||||
}
|
||||
|
||||
// ReaderToChan dumps all content from a given reader to a chan by constantly reading it until EOF.
|
||||
func ReaderToChan(stream chan<- *alloc.Buffer, reader Reader) error {
|
||||
for {
|
||||
buffer, err := reader.Read()
|
||||
if alloc.Len(buffer) > 0 {
|
||||
stream <- buffer
|
||||
} else {
|
||||
alloc.Release(buffer)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ChanToWriter dumps all content from a given chan to a writer until the chan is closed.
|
||||
func ChanToWriter(writer io.Writer, stream <-chan *alloc.Buffer) error {
|
||||
for buffer := range stream {
|
||||
nBytes, err := writer.Write(buffer.Value)
|
||||
if nBytes < buffer.Len() {
|
||||
_, err = writer.Write(buffer.Value[nBytes:])
|
||||
}
|
||||
buffer.Release()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
37
common/io/transport_test.go
Normal file
37
common/io/transport_test.go
Normal file
@@ -0,0 +1,37 @@
|
||||
package io_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto/rand"
|
||||
"io"
|
||||
"testing"
|
||||
|
||||
"github.com/v2ray/v2ray-core/common/alloc"
|
||||
. "github.com/v2ray/v2ray-core/common/io"
|
||||
v2testing "github.com/v2ray/v2ray-core/testing"
|
||||
"github.com/v2ray/v2ray-core/testing/assert"
|
||||
)
|
||||
|
||||
func TestReaderAndWrite(t *testing.T) {
|
||||
v2testing.Current(t)
|
||||
|
||||
size := 1024 * 1024
|
||||
buffer := make([]byte, size)
|
||||
nBytes, err := rand.Read(buffer)
|
||||
assert.Int(nBytes).Equals(len(buffer))
|
||||
assert.Error(err).IsNil()
|
||||
|
||||
readerBuffer := bytes.NewReader(buffer)
|
||||
writerBuffer := bytes.NewBuffer(make([]byte, 0, size))
|
||||
|
||||
transportChan := make(chan *alloc.Buffer, 1024)
|
||||
|
||||
err = ReaderToChan(transportChan, NewAdaptiveReader(readerBuffer))
|
||||
assert.Error(err).Equals(io.EOF)
|
||||
close(transportChan)
|
||||
|
||||
err = ChanToWriter(writerBuffer, transportChan)
|
||||
assert.Error(err).IsNil()
|
||||
|
||||
assert.Bytes(buffer).Equals(writerBuffer.Bytes())
|
||||
}
|
||||
@@ -1,96 +0,0 @@
|
||||
package net
|
||||
|
||||
import (
|
||||
"io"
|
||||
|
||||
"github.com/v2ray/v2ray-core/common/alloc"
|
||||
"github.com/v2ray/v2ray-core/common/crypto"
|
||||
"github.com/v2ray/v2ray-core/common/serial"
|
||||
"github.com/v2ray/v2ray-core/transport"
|
||||
)
|
||||
|
||||
// ReadFrom reads from a reader and put all content to a buffer.
|
||||
// If buffer is nil, ReadFrom creates a new normal buffer.
|
||||
func ReadFrom(reader io.Reader, buffer *alloc.Buffer) (*alloc.Buffer, error) {
|
||||
if buffer == nil {
|
||||
buffer = alloc.NewBuffer()
|
||||
}
|
||||
nBytes, err := reader.Read(buffer.Value)
|
||||
buffer.Slice(0, nBytes)
|
||||
return buffer, err
|
||||
}
|
||||
|
||||
func ReadChunk(reader io.Reader, buffer *alloc.Buffer) (*alloc.Buffer, error) {
|
||||
if buffer == nil {
|
||||
buffer = alloc.NewBuffer()
|
||||
}
|
||||
if _, err := io.ReadFull(reader, buffer.Value[:2]); err != nil {
|
||||
alloc.Release(buffer)
|
||||
return nil, err
|
||||
}
|
||||
length := serial.BytesLiteral(buffer.Value[:2]).Uint16Value()
|
||||
if _, err := io.ReadFull(reader, buffer.Value[:length]); err != nil {
|
||||
alloc.Release(buffer)
|
||||
return nil, err
|
||||
}
|
||||
buffer.Slice(0, int(length))
|
||||
return buffer, nil
|
||||
}
|
||||
|
||||
func ReadAuthenticatedChunk(reader io.Reader, auth crypto.Authenticator, buffer *alloc.Buffer) (*alloc.Buffer, error) {
|
||||
buffer, err := ReadChunk(reader, buffer)
|
||||
if err != nil {
|
||||
alloc.Release(buffer)
|
||||
return nil, err
|
||||
}
|
||||
authSize := auth.AuthBytes()
|
||||
|
||||
authBytes := auth.Authenticate(nil, buffer.Value[authSize:])
|
||||
|
||||
if !serial.BytesLiteral(authBytes).Equals(serial.BytesLiteral(buffer.Value[:authSize])) {
|
||||
alloc.Release(buffer)
|
||||
return nil, transport.CorruptedPacket
|
||||
}
|
||||
buffer.SliceFrom(authSize)
|
||||
|
||||
return buffer, nil
|
||||
}
|
||||
|
||||
// ReaderToChan dumps all content from a given reader to a chan by constantly reading it until EOF.
|
||||
func ReaderToChan(stream chan<- *alloc.Buffer, reader io.Reader) error {
|
||||
allocate := alloc.NewBuffer
|
||||
large := false
|
||||
for {
|
||||
buffer, err := ReadFrom(reader, allocate())
|
||||
if buffer.Len() > 0 {
|
||||
stream <- buffer
|
||||
} else {
|
||||
buffer.Release()
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if buffer.IsFull() && !large {
|
||||
allocate = alloc.NewLargeBuffer
|
||||
large = true
|
||||
} else if !buffer.IsFull() {
|
||||
allocate = alloc.NewBuffer
|
||||
large = false
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ChanToWriter dumps all content from a given chan to a writer until the chan is closed.
|
||||
func ChanToWriter(writer io.Writer, stream <-chan *alloc.Buffer) error {
|
||||
for buffer := range stream {
|
||||
nBytes, err := writer.Write(buffer.Value)
|
||||
if nBytes < buffer.Len() {
|
||||
_, err = writer.Write(buffer.Value[nBytes:])
|
||||
}
|
||||
buffer.Release()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -1,153 +0,0 @@
|
||||
package net_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto/rand"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"testing"
|
||||
|
||||
"github.com/v2ray/v2ray-core/common/alloc"
|
||||
v2net "github.com/v2ray/v2ray-core/common/net"
|
||||
v2testing "github.com/v2ray/v2ray-core/testing"
|
||||
"github.com/v2ray/v2ray-core/testing/assert"
|
||||
)
|
||||
|
||||
func TestReaderAndWrite(t *testing.T) {
|
||||
v2testing.Current(t)
|
||||
|
||||
size := 1024 * 1024
|
||||
buffer := make([]byte, size)
|
||||
nBytes, err := rand.Read(buffer)
|
||||
assert.Int(nBytes).Equals(len(buffer))
|
||||
assert.Error(err).IsNil()
|
||||
|
||||
readerBuffer := bytes.NewReader(buffer)
|
||||
writerBuffer := bytes.NewBuffer(make([]byte, 0, size))
|
||||
|
||||
transportChan := make(chan *alloc.Buffer, 1024)
|
||||
|
||||
err = v2net.ReaderToChan(transportChan, readerBuffer)
|
||||
assert.Error(err).Equals(io.EOF)
|
||||
close(transportChan)
|
||||
|
||||
err = v2net.ChanToWriter(writerBuffer, transportChan)
|
||||
assert.Error(err).IsNil()
|
||||
|
||||
assert.Bytes(buffer).Equals(writerBuffer.Bytes())
|
||||
}
|
||||
|
||||
type StaticReader struct {
|
||||
total int
|
||||
current int
|
||||
}
|
||||
|
||||
func (reader *StaticReader) Read(b []byte) (size int, err error) {
|
||||
size = len(b)
|
||||
if size > reader.total-reader.current {
|
||||
size = reader.total - reader.current
|
||||
}
|
||||
for i := 0; i < size; i++ {
|
||||
b[i] = byte(i)
|
||||
}
|
||||
//rand.Read(b[:size])
|
||||
reader.current += size
|
||||
if reader.current == reader.total {
|
||||
err = io.EOF
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func BenchmarkTransport1K(b *testing.B) {
|
||||
size := 1 * 1024
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
runBenchmarkTransport(size)
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkTransport2K(b *testing.B) {
|
||||
size := 2 * 1024
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
runBenchmarkTransport(size)
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkTransport4K(b *testing.B) {
|
||||
size := 4 * 1024
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
runBenchmarkTransport(size)
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkTransport10K(b *testing.B) {
|
||||
size := 10 * 1024
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
runBenchmarkTransport(size)
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkTransport100K(b *testing.B) {
|
||||
size := 100 * 1024
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
runBenchmarkTransport(size)
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkTransport1M(b *testing.B) {
|
||||
size := 1024 * 1024
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
runBenchmarkTransport(size)
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkTransport10M(b *testing.B) {
|
||||
size := 10 * 1024 * 1024
|
||||
|
||||
for i := 0; i < b.N; i++ {
|
||||
runBenchmarkTransport(size)
|
||||
}
|
||||
}
|
||||
|
||||
func runBenchmarkTransport(size int) {
|
||||
|
||||
transportChanA := make(chan *alloc.Buffer, 16)
|
||||
transportChanB := make(chan *alloc.Buffer, 16)
|
||||
|
||||
readerA := &StaticReader{size, 0}
|
||||
readerB := &StaticReader{size, 0}
|
||||
|
||||
writerA := ioutil.Discard
|
||||
writerB := ioutil.Discard
|
||||
|
||||
finishA := make(chan bool)
|
||||
finishB := make(chan bool)
|
||||
|
||||
go func() {
|
||||
v2net.ChanToWriter(writerA, transportChanA)
|
||||
close(finishA)
|
||||
}()
|
||||
|
||||
go func() {
|
||||
v2net.ReaderToChan(transportChanA, readerA)
|
||||
close(transportChanA)
|
||||
}()
|
||||
|
||||
go func() {
|
||||
v2net.ChanToWriter(writerB, transportChanB)
|
||||
close(finishB)
|
||||
}()
|
||||
|
||||
go func() {
|
||||
v2net.ReaderToChan(transportChanB, readerB)
|
||||
close(transportChanB)
|
||||
}()
|
||||
|
||||
<-transportChanA
|
||||
<-transportChanB
|
||||
}
|
||||
Reference in New Issue
Block a user