From 33e0509c5a5ba044b2fea9a70fdafe9079ce3e75 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Wed, 1 Dec 2021 10:51:13 +0300 Subject: [PATCH] Optimize for a fast flush --- mtglib/internal/relay/init.go | 4 ++-- mtglib/internal/relay/sync_pair.go | 13 ++++++++++++- 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/mtglib/internal/relay/init.go b/mtglib/internal/relay/init.go index e54ac7a..6278c48 100644 --- a/mtglib/internal/relay/init.go +++ b/mtglib/internal/relay/init.go @@ -3,8 +3,8 @@ package relay import "time" const ( - copyBufferSize = 32 * 1024 - writerBufferSize = 2 * copyBufferSize + copyBufferSize = 64 * 1024 + writerBufferSize = 128 * 1024 readTimeout = 10 * time.Millisecond ) diff --git a/mtglib/internal/relay/sync_pair.go b/mtglib/internal/relay/sync_pair.go index 72293e7..500b961 100644 --- a/mtglib/internal/relay/sync_pair.go +++ b/mtglib/internal/relay/sync_pair.go @@ -26,6 +26,7 @@ func (s *syncPair) Sync() (int64, error) { func (s *syncPair) Read(p []byte) (int, error) { n, err := s.readBlocking(p, false) + // nothing has been delivered for readTimeout time. Let's flush. if errors.Is(err, os.ErrDeadlineExceeded) { if err := s.Flush(); err != nil { return 0, fmt.Errorf("cannot flush writer hand-side: %w", err) @@ -41,7 +42,17 @@ func (s *syncPair) Write(p []byte) (int, error) { s.mutex.Lock() defer s.mutex.Unlock() - return s.writer.Write(p) // nolint: wrapcheck + n, err := s.writer.Write(p) // nolint: wrapcheck + + // optimization for a case when we have a small package and want to avoid a + // delay in readTimeout. In that case, we assume that peer has finished to + // sent a data it wants to send so we can flush without waiting for anything + // else. + if err == nil && n < copyBufferSize { + err = s.writer.Flush() + } + + return n, err } func (s *syncPair) Flush() error {