pref(rpc2): _gc_ reuse objects in transport frames

. #23113571

cpu: Intel(R) Core(TM) i7-10700 CPU @ 2.90GHz
BenchmarkAcceptClose
BenchmarkAcceptClose  276475   5393 ns/op  904 B/op  7 allocs/op
BenchmarkConnSmux
BenchmarkConnSmux      39103  31170 ns/op  121 B/op  0 allocs/op
BenchmarkPipeSmux
BenchmarkPipeSmux      90400  13771 ns/op    6 B/op  0 allocs/op
BenchmarkConnTCP
BenchmarkConnTCP       61194  18310 ns/op    0 B/op  0 allocs/op
BenchmarkRangedWrite
BenchmarkRangedWrite  242226   5230 ns/op   24 B/op  0 allocs/op

Signed-off-by: slasher <shenjie1@oppo.com>
This commit is contained in:
slasher 2025-03-06 15:47:10 +08:00
parent a793d19ce0
commit e5125b59c4
6 changed files with 189 additions and 87 deletions

View File

@ -68,11 +68,14 @@ type netConnv struct{ netConn }
var _ WriteBuffers = netConnv{}
func (c netConnv) WriteBuffers(buffers []AssignedBuffer) (int, error) {
v := make(net.Buffers, 0, len(buffers))
pv := poolBuffers.Get().(*net.Buffers)
v := (*pv)[:0]
for _, buffer := range buffers {
v = append(v, buffer.Bytes()[:buffer.Len()])
}
*pv = v
nn, err := v.WriteTo(c.netConn.Conn)
poolBuffers.Put(pv) // nolint: staticcheck
return int(nn), err
}

View File

@ -83,6 +83,8 @@ func (h updHeader) Window() uint32 {
// FrameWrite frame for write
type FrameWrite struct {
recycle bool
ver byte
cmd byte
sid uint32
@ -147,7 +149,19 @@ func (f *FrameWrite) Len() int {
func (f *FrameWrite) Close() (err error) {
if f.f != nil {
return f.f.Close()
if f.recycle {
f.f.recycle = true
}
err = f.f.Close()
if f.recycle {
if ab, is := f.ab.(*unAlignedBuffer); is {
*ab = unAlignedBuffer{}
poolunAlignedBuffer.Put(ab) // nolint: staticcheck
}
*f = FrameWrite{}
poolFrameWrite.Put(f) // nolint: staticcheck
}
return
}
for !atomic.CompareAndSwapUint32(&f.done, 0, 2) {
if atomic.LoadUint32(&f.done) == 2 {
@ -159,6 +173,10 @@ func (f *FrameWrite) Close() (err error) {
if f.ab != nil {
err = f.ab.Free()
f.ab = nil
if f.recycle {
*f = FrameWrite{}
poolFrameWrite.Put(f) // nolint: staticcheck
}
}
return
}
@ -188,12 +206,14 @@ func (ub *unAlignedBuffer) Free() error { return ub.ab.Free() }
func (f *FrameWrite) TrimHead(head int) *FrameWrite {
off := f.off - head
data := f.data[head:]
ab := &unAlignedBuffer{
ab := poolunAlignedBuffer.Get().(*unAlignedBuffer)
*ab = unAlignedBuffer{
offset: off,
buffer: data,
ab: f.ab,
}
return &FrameWrite{
fw := poolFrameWrite.Get().(*FrameWrite)
*fw = FrameWrite{
ver: f.ver,
cmd: f.cmd,
sid: f.sid,
@ -205,6 +225,7 @@ func (f *FrameWrite) TrimHead(head int) *FrameWrite {
ctx: f.ctx,
f: f,
}
return fw
}
func (f *FrameWrite) TrimTail(tail int) *FrameWrite {
@ -263,6 +284,7 @@ func (f *FrameRead) Close() (err error) {
if f.ab != nil {
err = f.ab.Free()
f.ab = nil
poolFrameRead.Put(f) // nolint: staticcheck
}
return
}

View File

@ -46,11 +46,6 @@ type writeResult struct {
err error
}
type writeDealine struct {
time time.Time
wait <-chan time.Time
}
// Session defines a multiplexed connection for streams
type Session struct {
conn Conn
@ -174,8 +169,8 @@ func (s *Session) OpenStream() (*Stream, error) {
func (s *Session) AcceptStream() (*Stream, error) {
var deadline <-chan time.Time
if d, ok := s.deadline.Load().(time.Time); ok && !d.IsZero() {
timer := time.NewTimer(time.Until(d))
defer timer.Stop()
timer := acquirePoolTimer(time.Until(d))
defer releasePoolTimer(timer)
deadline = timer.C
}
@ -392,13 +387,10 @@ func (s *Session) recvLoop() {
}
func (s *Session) ping(ticker *time.Ticker) {
deadline := writeDealine{
time: time.Now().Add(s.config.KeepAliveInterval),
wait: ticker.C,
}
dur := s.config.KeepAliveInterval
frame, err := s.newFrameWrite(cmdPIN, 0, 0)
if err == nil {
s.writeFrameInternal(frame, deadline, CLSCTRL)
s.writeFrameInternal(frame, CLSCTRL, time.Now().Add(dur), ticker.C)
frame.Close()
}
}
@ -562,21 +554,29 @@ func (s *Session) sendLoop() {
// writeFrame writes the frame to the underlying connection
// and returns the number of bytes written if successful
func (s *Session) writeFrame(f *FrameWrite) (n int, err error) {
timer := time.NewTimer(openCloseTimeout)
defer timer.Stop()
defer f.Close()
deadline := writeDealine{
time: time.Now().Add(openCloseTimeout),
wait: timer.C,
}
return s.writeFrameInternal(f, deadline, CLSCTRL)
n, err = s.writeFrameInternal(f, CLSCTRL, time.Now().Add(openCloseTimeout), nil)
f.Close()
return
}
// internal writeFrame version to support deadline used in keepalive
func (s *Session) writeFrameInternal(f *FrameWrite, deadline writeDealine, class CLASSID) (int, error) {
func (s *Session) writeFrameInternal(f *FrameWrite, class CLASSID,
deadline time.Time, deadwait <-chan time.Time,
) (int, error) {
var timer *time.Timer
if !deadline.IsZero() && deadwait == nil {
tu := time.Until(deadline)
if tu <= 0 {
return 0, ErrTimeout
}
timer = acquirePoolTimer(tu)
defer releasePoolTimer(timer)
deadwait = timer.C
}
req := writeRequest{
frame: f,
deadline: deadline.time,
deadline: deadline,
result: s.resultChPool.Get().(chan writeResult),
}
writeCh := s.writes
@ -586,7 +586,7 @@ func (s *Session) writeFrameInternal(f *FrameWrite, deadline writeDealine, class
ctx := f.Context()
select {
case <-deadline.wait:
case <-deadwait:
return 0, ErrTimeout
default:
select {
@ -595,7 +595,7 @@ func (s *Session) writeFrameInternal(f *FrameWrite, deadline writeDealine, class
return 0, io.ErrClosedPipe
case <-s.chSocketWriteError:
return 0, s.socketWriteError.Load().(error)
case <-deadline.wait:
case <-deadwait:
return 0, ErrTimeout
case <-ctx.Done():
return 0, ctx.Err()
@ -605,6 +605,7 @@ func (s *Session) writeFrameInternal(f *FrameWrite, deadline writeDealine, class
select {
case result := <-req.result:
s.resultChPool.Put(req.result)
f.recycle = true
return result.n, result.err
case <-s.die:
return 0, io.ErrClosedPipe
@ -621,7 +622,8 @@ func (s *Session) newFrameWrite(cmd byte, sid uint32, size int) (*FrameWrite, er
return nil, err
}
buffer.Written(headerSize)
return &FrameWrite{
fw := poolFrameWrite.Get().(*FrameWrite)
*fw = FrameWrite{
ver: byte(s.config.Version),
cmd: cmd,
sid: sid,
@ -629,14 +631,17 @@ func (s *Session) newFrameWrite(cmd byte, sid uint32, size int) (*FrameWrite, er
ab: buffer,
off: headerSize,
data: buffer.Bytes()[:],
}, nil
}
return fw, nil
}
func (s *Session) newFrameRead(buffer AssignedBuffer) *FrameRead {
return &FrameRead{
fr := poolFrameRead.Get().(*FrameRead)
*fr = FrameRead{
ab: buffer,
data: buffer.Bytes()[headerSize:buffer.Len()],
}
return fr
}
func longestTime(t time.Time, other time.Time) time.Time {

View File

@ -1093,11 +1093,9 @@ func TestWriteFrameInternal(t *testing.T) {
session.Close()
for i := 0; i < 100; i++ {
f, _ := session.newFrameWrite(byte(rand.Uint32()), rand.Uint32(), 0)
deadline := writeDealine{
time: time.Now().Add(session.config.KeepAliveTimeout),
wait: time.After(session.config.KeepAliveTimeout),
}
session.writeFrameInternal(f, deadline, CLSDATA)
session.writeFrameInternal(f, CLSDATA,
time.Now().Add(session.config.KeepAliveTimeout),
time.After(session.config.KeepAliveTimeout))
}
// random cmds
@ -1109,18 +1107,16 @@ func TestWriteFrameInternal(t *testing.T) {
session, _ = Client(newConn(cli), nil)
for i := 0; i < 100; i++ {
f, _ := session.newFrameWrite(allcmds[rand.Int()%len(allcmds)], rand.Uint32(), 0)
deadline := writeDealine{
time: time.Now().Add(session.config.KeepAliveTimeout),
wait: time.After(session.config.KeepAliveTimeout),
}
session.writeFrameInternal(f, deadline, CLSDATA)
session.writeFrameInternal(f, CLSDATA,
time.Now().Add(session.config.KeepAliveTimeout),
time.After(session.config.KeepAliveTimeout))
}
// deadline occur
{
c := make(chan time.Time)
close(c)
f, _ := session.newFrameWrite(allcmds[rand.Int()%len(allcmds)], rand.Uint32(), 0)
_, err = session.writeFrameInternal(f, writeDealine{wait: c}, CLSDATA)
_, err = session.writeFrameInternal(f, CLSDATA, time.Time{}, c)
if !strings.Contains(err.Error(), "timeout") {
t.Fatal("write frame with deadline failed", err)
}
@ -1145,7 +1141,7 @@ func TestWriteFrameInternal(t *testing.T) {
time.Sleep(time.Second)
close(c)
}()
_, err = session.writeFrameInternal(f, writeDealine{wait: c}, CLSDATA)
_, err = session.writeFrameInternal(f, CLSDATA, time.Time{}, c)
if !strings.Contains(err.Error(), "closed pipe") {
t.Fatal("write frame with to closed conn failed", err)
}
@ -1437,3 +1433,48 @@ func bench(b *testing.B, rd io.Reader, wr io.Writer) {
}
wg.Wait()
}
func BenchmarkRangedWrite(b *testing.B) {
cs, ss, err := getSmuxStreamPair(false)
if err != nil {
b.Fatal(err)
}
defer cs.Close()
defer ss.Close()
const size = (2 * 1024)
buf := make([]byte, size)
cr := bytes.NewReader(buf)
b.SetBytes(size)
b.ResetTimer()
b.ReportAllocs()
b.SetParallelism(1)
runtime.GC()
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
count := 0
for {
ss.SetReadDeadline(time.Now().Add(time.Minute))
sr := ss.NewSizedReader(testCtx, size-20, nil)
n, err := sr.WriteTo(io.Discard)
sr.Close()
if err != nil {
panic(err)
}
count += int(n)
if count == (size-20)*b.N {
return
}
}
}()
for i := 0; i < b.N; i++ {
cr.Reset(buf)
cs.SetWriteDeadline(time.Now().Add(time.Minute))
cs.RangedWrite(testCtx, cr, size, 10, 10, false, nil)
}
wg.Wait()
}

View File

@ -31,6 +31,8 @@ type Stream struct {
finEventOnce sync.Once
// deadlines
readTime time.Time
writeTime time.Time
readDeadline atomic.Value
writeDeadline atomic.Value
@ -160,16 +162,9 @@ func (s *Stream) tryReadFramev2() (*FrameRead, error) {
}
func (s *Stream) sendWindowUpdate(consumed uint32) error {
var deadline writeDealine
if d, ok := s.readDeadline.Load().(time.Time); ok && !d.IsZero() {
tu := time.Until(d)
if tu <= 0 {
return ErrTimeout
}
timer := time.NewTimer(tu)
defer timer.Stop()
deadline.time = d
deadline.wait = timer.C
var deadline time.Time
if d, ok := s.readDeadline.Load().(*time.Time); ok && !(*d).IsZero() {
deadline = *d
}
var hdr updHeader
@ -180,16 +175,16 @@ func (s *Stream) sendWindowUpdate(consumed uint32) error {
binary.LittleEndian.PutUint32(hdr[:], consumed)
binary.LittleEndian.PutUint32(hdr[4:], uint32(s.sess.config.MaxStreamBuffer))
frame.Write(hdr[:])
_, err = s.sess.writeFrameInternal(frame, deadline, CLSDATA)
_, err = s.sess.writeFrameInternal(frame, CLSDATA, deadline, nil)
return err
}
func (s *Stream) waitRead(ctx context.Context) error {
var timer *time.Timer
var deadline <-chan time.Time
if d, ok := s.readDeadline.Load().(time.Time); ok && !d.IsZero() {
timer = time.NewTimer(time.Until(d))
defer timer.Stop()
if d, ok := s.readDeadline.Load().(*time.Time); ok && !(*d).IsZero() {
timer = acquirePoolTimer(time.Until(*d))
defer releasePoolTimer(timer)
deadline = timer.C
}
@ -316,28 +311,32 @@ func (s *Stream) WriteFrame(frame *FrameWrite) (n int, err error) {
default:
}
var deadline writeDealine
if d, ok := s.writeDeadline.Load().(time.Time); ok && !d.IsZero() {
tu := time.Until(d)
if tu <= 0 {
return 0, ErrTimeout
}
timer := time.NewTimer(tu)
defer timer.Stop()
deadline.time = d
deadline.wait = timer.C
var deadline time.Time
if d, ok := s.writeDeadline.Load().(*time.Time); ok && !(*d).IsZero() {
deadline = *d
}
if s.sess.config.Version == 2 {
return s.writeFrameV2(frame, deadline)
}
sent, err := s.sess.writeFrameInternal(frame, deadline, CLSDATA)
sent, err := s.sess.writeFrameInternal(frame, CLSDATA, deadline, nil)
s.numWritten += uint32(sent)
return sent, err
}
func (s *Stream) writeFrameV2(frame *FrameWrite, deadline writeDealine) (n int, err error) {
func (s *Stream) writeFrameV2(frame *FrameWrite, deadline time.Time) (n int, err error) {
var timer *time.Timer
var deadwait <-chan time.Time
if !deadline.IsZero() {
tu := time.Until(deadline)
if tu <= 0 {
return 0, ErrTimeout
}
timer = acquirePoolTimer(tu)
defer releasePoolTimer(timer)
deadwait = timer.C
}
for {
// per stream sliding window control
// [.... [consumed... numWritten] ... win... ]
@ -355,7 +354,7 @@ func (s *Stream) writeFrameV2(frame *FrameWrite, deadline writeDealine) (n int,
win := int32(atomic.LoadUint32(&s.peerWindow)) - inflight
if win >= int32(frame.Len()) || s.numWritten == 0 {
sent, err := s.sess.writeFrameInternal(frame, deadline, CLSDATA)
sent, err := s.sess.writeFrameInternal(frame, CLSDATA, deadline, deadwait)
s.numWritten += uint32(sent)
return sent, err
}
@ -367,7 +366,7 @@ func (s *Stream) writeFrameV2(frame *FrameWrite, deadline writeDealine) (n int,
return 0, io.EOF
case <-s.die:
return 0, io.ErrClosedPipe
case <-deadline.wait:
case <-deadwait:
return 0, ErrTimeout
case <-s.sess.chSocketWriteError:
return 0, s.sess.socketWriteError.Load().(error)
@ -421,7 +420,8 @@ func (s *Stream) IsClosed() bool {
// net.Conn.SetReadDeadline.
// A zero time value disables the deadline.
func (s *Stream) SetReadDeadline(t time.Time) error {
s.readDeadline.Store(t)
s.readTime = t
s.readDeadline.Store(&s.readTime)
s.notifyReadEvent()
return nil
}
@ -430,7 +430,8 @@ func (s *Stream) SetReadDeadline(t time.Time) error {
// net.Conn.SetWriteDeadline.
// A zero time value disables the deadline.
func (s *Stream) SetWriteDeadline(t time.Time) error {
s.writeDeadline.Store(t)
s.writeTime = t
s.writeDeadline.Store(&s.writeTime)
return nil
}
@ -502,13 +503,13 @@ type SizedReader struct {
f *FrameRead
finished bool
once sync.Once
err error
err error
}
func (s *Stream) NewSizedReader(ctx context.Context, size int, f *FrameRead) *SizedReader {
return &SizedReader{ctx: ctx, n: size, s: s, f: f}
sz := poolSizedReader.Get().(*SizedReader)
*sz = SizedReader{ctx: ctx, n: size, s: s, f: f}
return sz
}
func (r *SizedReader) tryNextFrame() error {
@ -520,6 +521,7 @@ func (r *SizedReader) tryNextFrame() error {
r.err = io.EOF
if r.f != nil && r.f.Len() == 0 {
r.f.Close()
r.f = nil
}
return r.err
}
@ -567,6 +569,9 @@ func (r *SizedReader) WriteTo(w io.Writer) (int64, error) {
return nn, nil
}
if err != nil {
if nn > 0 && err == io.EOF {
err = nil
}
return nn, err
}
}
@ -574,16 +579,10 @@ func (r *SizedReader) WriteTo(w io.Writer) (int64, error) {
func (r *SizedReader) Close() (err error) {
r.tryNextFrame()
r.once.Do(func() {
if r.err == nil {
r.err = io.EOF
}
if r.f != nil {
r.f.Close()
}
r.s = nil
r.f = nil
})
if r.s != nil {
*r = SizedReader{}
poolSizedReader.Put(r) // nolint: staticcheck
}
return
}

View File

@ -8,6 +8,8 @@ import (
"errors"
"fmt"
"math"
"net"
"sync"
"time"
)
@ -106,3 +108,33 @@ func Client(conn Conn, config *Config) (*Session, error) {
}
return newSession(config, conn, true), nil
}
var (
poolFrameWrite = sync.Pool{New: func() any { return new(FrameWrite) }}
poolFrameRead = sync.Pool{New: func() any { return new(FrameRead) }}
poolSizedReader = sync.Pool{New: func() any { return new(SizedReader) }}
poolunAlignedBuffer = sync.Pool{New: func() any { return new(unAlignedBuffer) }}
poolTimer = sync.Pool{New: func() any { return time.NewTimer(time.Hour) }}
poolBuffers = sync.Pool{New: func() any { b := make(net.Buffers, 32); return &b }}
)
func acquirePoolTimer(d time.Duration) *time.Timer {
timer := poolTimer.Get().(*time.Timer)
timer.Stop()
select {
case <-timer.C:
default:
}
timer.Reset(d)
return timer
}
func releasePoolTimer(timer *time.Timer) {
timer.Stop()
select {
case <-timer.C:
default:
}
poolTimer.Put(timer) // nolint:staticcheck
}