1
0
mirror of https://github.com/v2fly/v2ray-core.git synced 2024-12-22 01:57:12 -05:00

reuse buffered writer in auth writer

This commit is contained in:
Darien Raymond 2017-05-24 00:54:30 +02:00
parent d3e262d474
commit ade88fd5c7
No known key found for this signature in database
GPG Key ID: 7251FFA14BB18169
5 changed files with 63 additions and 56 deletions

View File

@ -11,10 +11,14 @@ type BufferedWriter struct {
} }
// NewBufferedWriter creates a new BufferedWriter. // NewBufferedWriter creates a new BufferedWriter.
func NewBufferedWriter(rawWriter io.Writer) *BufferedWriter { func NewBufferedWriter(writer io.Writer) *BufferedWriter {
return NewBufferedWriterSize(writer, 1024)
}
func NewBufferedWriterSize(writer io.Writer, size uint32) *BufferedWriter {
return &BufferedWriter{ return &BufferedWriter{
writer: rawWriter, writer: writer,
buffer: NewLocal(1024), buffer: NewLocal(int(size)),
buffered: true, buffered: true,
} }
} }
@ -24,21 +28,20 @@ func (w *BufferedWriter) Write(b []byte) (int, error) {
if !w.buffered || w.buffer == nil { if !w.buffered || w.buffer == nil {
return w.writer.Write(b) return w.writer.Write(b)
} }
nBytes, err := w.buffer.Write(b) bytesWritten := 0
for bytesWritten < len(b) {
nBytes, err := w.buffer.Write(b[bytesWritten:])
if err != nil { if err != nil {
return 0, err return bytesWritten, err
} }
bytesWritten += nBytes
if w.buffer.IsFull() { if w.buffer.IsFull() {
if err := w.Flush(); err != nil { if err := w.Flush(); err != nil {
return 0, err return bytesWritten, err
}
if nBytes < len(b) {
if _, err := w.writer.Write(b[nBytes:]); err != nil {
return nBytes, err
} }
} }
} }
return len(b), nil return bytesWritten, nil
} }
// Flush writes all buffered content into underlying writer, if any. // Flush writes all buffered content into underlying writer, if any.

View File

@ -47,6 +47,7 @@ func TestBufferedWriterLargePayload(t *testing.T) {
nBytes, err = writer.Write(payload[512:]) nBytes, err = writer.Write(payload[512:])
assert.Error(err).IsNil() assert.Error(err).IsNil()
assert.Error(writer.Flush()).IsNil()
assert.Int(nBytes).Equals(64*1024 - 512) assert.Int(nBytes).Equals(64*1024 - 512)
assert.Bytes(content.Bytes()).Equals(payload) assert.Bytes(content.Bytes()).Equals(payload)
} }

View File

@ -4,6 +4,7 @@ import (
"crypto/cipher" "crypto/cipher"
"io" "io"
"v2ray.com/core/common"
"v2ray.com/core/common/buf" "v2ray.com/core/common/buf"
"v2ray.com/core/common/protocol" "v2ray.com/core/common/protocol"
) )
@ -123,7 +124,12 @@ func (r *AuthenticationReader) readChunk(waitForData bool) ([]byte, error) {
if !waitForData { if !waitForData {
return nil, io.ErrNoProgress return nil, io.ErrNoProgress
} }
r.buffer.Reset(buf.ReadFrom(r.buffer))
if r.buffer.IsEmpty() {
r.buffer.Clear()
} else {
common.Must(r.buffer.Reset(buf.ReadFrom(r.buffer)))
}
delta := r.size - r.buffer.Len() delta := r.size - r.buffer.Len()
if err := r.buffer.AppendSupplier(buf.ReadAtLeastFrom(r.reader, delta)); err != nil { if err := r.buffer.AppendSupplier(buf.ReadAtLeastFrom(r.reader, delta)); err != nil {
@ -184,64 +190,55 @@ func (r *AuthenticationReader) Read() (buf.MultiBuffer, error) {
type AuthenticationWriter struct { type AuthenticationWriter struct {
auth Authenticator auth Authenticator
buffer []byte
payload []byte payload []byte
buffer *buf.Buffer writer *buf.BufferedWriter
writer io.Writer
sizeParser ChunkSizeEncoder sizeParser ChunkSizeEncoder
transferType protocol.TransferType transferType protocol.TransferType
} }
func NewAuthenticationWriter(auth Authenticator, sizeParser ChunkSizeEncoder, writer io.Writer, transferType protocol.TransferType) *AuthenticationWriter { func NewAuthenticationWriter(auth Authenticator, sizeParser ChunkSizeEncoder, writer io.Writer, transferType protocol.TransferType) *AuthenticationWriter {
const payloadSize = 1024
return &AuthenticationWriter{ return &AuthenticationWriter{
auth: auth, auth: auth,
payload: make([]byte, 1024), buffer: make([]byte, payloadSize+sizeParser.SizeBytes()+auth.Overhead()),
buffer: buf.NewLocal(readerBufferSize), payload: make([]byte, payloadSize),
writer: writer, writer: buf.NewBufferedWriterSize(writer, readerBufferSize),
sizeParser: sizeParser, sizeParser: sizeParser,
transferType: transferType, transferType: transferType,
} }
} }
func (w *AuthenticationWriter) append(b []byte) { func (w *AuthenticationWriter) append(b []byte) error {
encryptedSize := len(b) + w.auth.Overhead() encryptedSize := len(b) + w.auth.Overhead()
buffer := w.sizeParser.Encode(uint16(encryptedSize), w.buffer[:0])
w.buffer.AppendSupplier(func(bb []byte) (int, error) { buffer, err := w.auth.Seal(buffer, b)
w.sizeParser.Encode(uint16(encryptedSize), bb[:0]) if err != nil {
return w.sizeParser.SizeBytes(), nil return err
})
w.buffer.AppendSupplier(func(bb []byte) (int, error) {
w.auth.Seal(bb[:0], b)
return encryptedSize, nil
})
} }
func (w *AuthenticationWriter) flush() error { if _, err := w.writer.Write(buffer); err != nil {
_, err := w.writer.Write(w.buffer.Bytes())
w.buffer.Clear()
return err return err
} }
return nil
}
func (w *AuthenticationWriter) writeStream(mb buf.MultiBuffer) error { func (w *AuthenticationWriter) writeStream(mb buf.MultiBuffer) error {
defer mb.Release() defer mb.Release()
for { for {
n, _ := mb.Read(w.payload) n, _ := mb.Read(w.payload)
w.append(w.payload[:n]) if err := w.append(w.payload[:n]); err != nil {
if w.buffer.Len() > readerBufferSize-2*1024 {
if err := w.flush(); err != nil {
return err return err
} }
}
if mb.IsEmpty() { if mb.IsEmpty() {
break break
} }
} }
if !w.buffer.IsEmpty() { return w.writer.Flush()
return w.flush()
}
return nil
} }
func (w *AuthenticationWriter) writePacket(mb buf.MultiBuffer) error { func (w *AuthenticationWriter) writePacket(mb buf.MultiBuffer) error {
@ -252,24 +249,17 @@ func (w *AuthenticationWriter) writePacket(mb buf.MultiBuffer) error {
if b == nil { if b == nil {
b = buf.New() b = buf.New()
} }
if w.buffer.Len() > readerBufferSize-b.Len()-128 { if err := w.append(b.Bytes()); err != nil {
if err := w.flush(); err != nil {
b.Release() b.Release()
return err return err
} }
}
w.append(b.Bytes())
b.Release() b.Release()
if mb.IsEmpty() { if mb.IsEmpty() {
break break
} }
} }
if !w.buffer.IsEmpty() { return w.writer.Flush()
return w.flush()
}
return nil
} }
func (w *AuthenticationWriter) Write(mb buf.MultiBuffer) error { func (w *AuthenticationWriter) Write(mb buf.MultiBuffer) error {

View File

@ -181,7 +181,7 @@ func (v *Handler) Process(ctx context.Context, network net.Network, connection i
if err != nil { if err != nil {
if errors.Cause(err) != io.EOF { if errors.Cause(err) != io.EOF {
log.Access(connection.RemoteAddr(), "", log.AccessRejected, err) log.Access(connection.RemoteAddr(), "", log.AccessRejected, err)
log.Trace(newError("invalid request from ", connection.RemoteAddr(), ": ", err)) log.Trace(newError("invalid request from ", connection.RemoteAddr(), ": ", err).AtInfo())
} }
return err return err
} }
@ -194,7 +194,7 @@ func (v *Handler) Process(ctx context.Context, network net.Network, connection i
log.Access(connection.RemoteAddr(), request.Destination(), log.AccessAccepted, "") log.Access(connection.RemoteAddr(), request.Destination(), log.AccessAccepted, "")
log.Trace(newError("received request for ", request.Destination())) log.Trace(newError("received request for ", request.Destination()))
connection.SetReadDeadline(time.Time{}) common.Must(connection.SetReadDeadline(time.Time{}))
userSettings := request.User.GetSettings() userSettings := request.User.GetSettings()

View File

@ -5,6 +5,7 @@ import (
"testing" "testing"
"v2ray.com/core" "v2ray.com/core"
"v2ray.com/core/app/log"
"v2ray.com/core/app/proxyman" "v2ray.com/core/app/proxyman"
v2net "v2ray.com/core/common/net" v2net "v2ray.com/core/common/net"
"v2ray.com/core/common/protocol" "v2ray.com/core/common/protocol"
@ -55,6 +56,12 @@ func TestDokodemoTCP(t *testing.T) {
ProxySettings: serial.ToTypedMessage(&freedom.Config{}), ProxySettings: serial.ToTypedMessage(&freedom.Config{}),
}, },
}, },
App: []*serial.TypedMessage{
serial.ToTypedMessage(&log.Config{
ErrorLogLevel: log.LogLevel_Debug,
ErrorLogType: log.LogType_Console,
}),
},
} }
clientPort := uint32(pickPort()) clientPort := uint32(pickPort())
@ -94,6 +101,12 @@ func TestDokodemoTCP(t *testing.T) {
}), }),
}, },
}, },
App: []*serial.TypedMessage{
serial.ToTypedMessage(&log.Config{
ErrorLogLevel: log.LogLevel_Debug,
ErrorLogType: log.LogType_Console,
}),
},
} }
servers, err := InitializeServerConfigs(serverConfig, clientConfig) servers, err := InitializeServerConfigs(serverConfig, clientConfig)