From 913a38d13a3ad1cd34e95d7489f808141bffb4e0 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Wed, 18 Mar 2026 22:04:02 +0100 Subject: [PATCH 1/6] Show real IP of the telegram endpoint in event stream --- mtglib/proxy.go | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/mtglib/proxy.go b/mtglib/proxy.go index 7eb31a6..0e66e7c 100644 --- a/mtglib/proxy.go +++ b/mtglib/proxy.go @@ -259,9 +259,16 @@ func (p *Proxy) doTelegramCall(ctx *streamContext) error { ctx: ctx, } + telegramHost, _, err := net.SplitHostPort(foundAddr.Address) + if err != nil { + conn.Close() //nolint: errcheck + + return fmt.Errorf("cannot parse telegram address %s: %w", foundAddr.Address, err) + } + p.eventStream.Send(ctx, NewEventConnectedToDC(ctx.streamID, - conn.RemoteAddr().(*net.TCPAddr).IP, //nolint: forcetypeassert + net.ParseIP(telegramHost), ctx.dc), ) From a23ae05f3b15c411a917330011607de8a8774344 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Thu, 19 Mar 2026 13:47:08 +0100 Subject: [PATCH 2/6] Remove SyncWrite --- mtglib/internal/doppel/conn.go | 26 ------ mtglib/internal/doppel/conn_test.go | 130 ---------------------------- 2 files changed, 156 deletions(-) diff --git a/mtglib/internal/doppel/conn.go b/mtglib/internal/doppel/conn.go index 7e8ed30..bd76c7f 100644 --- a/mtglib/internal/doppel/conn.go +++ b/mtglib/internal/doppel/conn.go @@ -36,32 +36,6 @@ func (c Conn) Write(p []byte) (int, error) { return len(p), context.Cause(c.p.ctx) } -func (c Conn) SyncWrite(p []byte) (int, error) { - c.p.syncWriteLock.Lock() - defer c.p.syncWriteLock.Unlock() - - c.p.writeCond.L.Lock() - // wait until buffer is exhausted - for c.p.writeStream.Len() != 0 && context.Cause(c.p.ctx) == nil { - c.p.writeCond.Wait() - } - c.p.writeStream.Write(p) - c.p.writeCond.L.Unlock() - - if err := context.Cause(c.p.ctx); err != nil { - return len(p), err - } - - c.p.writeCond.L.Lock() - // wait until data will be sent - for c.p.writeStream.Len() != 0 && context.Cause(c.p.ctx) == nil { - c.p.writeCond.Wait() - } - c.p.writeCond.L.Unlock() - - return len(p), context.Cause(c.p.ctx) -} - func (c Conn) Start() { c.p.wg.Go(func() { c.start() diff --git a/mtglib/internal/doppel/conn_test.go b/mtglib/internal/doppel/conn_test.go index eec1f6f..e774bce 100644 --- a/mtglib/internal/doppel/conn_test.go +++ b/mtglib/internal/doppel/conn_test.go @@ -157,136 +157,6 @@ func (suite *ConnTestSuite) TestStopOnUnderlyingWriteError() { }, 2*time.Second, time.Millisecond) } -func (suite *ConnTestSuite) TestSyncWriteDataSent() { - suite.connMock. - On("Write", mock.AnythingOfType("[]uint8")). - Return(0, nil). - Maybe() - - c := suite.makeConn() - defer c.Stop() - - payload := []byte("sync hello") - n, err := c.SyncWrite(payload) - suite.NoError(err) - suite.Equal(len(payload), n) - - // SyncWrite returns only after data is flushed to the wire. - assembled := &bytes.Buffer{} - reader := bytes.NewReader(suite.connMock.Written()) - - for { - header := make([]byte, tls.SizeHeader) - if _, err := io.ReadFull(reader, header); err != nil { - break - } - - suite.Equal(byte(tls.TypeApplicationData), header[0]) - - length := binary.BigEndian.Uint16(header[tls.SizeRecordType+tls.SizeVersion:]) - rec := make([]byte, length) - _, err := io.ReadFull(reader, rec) - suite.NoError(err) - - assembled.Write(rec) - } - - suite.Equal(payload, assembled.Bytes()) -} - -func (suite *ConnTestSuite) TestSyncWriteDrainsBufferFirst() { - suite.connMock. - On("Write", mock.AnythingOfType("[]uint8")). - Return(0, nil). - Maybe() - - c := suite.makeConn() - defer c.Stop() - - // Buffer some data via async Write. - _, err := c.Write([]byte("first")) - suite.NoError(err) - - // SyncWrite must drain "first" before sending "second". - n, err := c.SyncWrite([]byte("second")) - suite.NoError(err) - suite.Equal(6, n) - - // All data should be on the wire now. - assembled := &bytes.Buffer{} - reader := bytes.NewReader(suite.connMock.Written()) - - for { - header := make([]byte, tls.SizeHeader) - if _, err := io.ReadFull(reader, header); err != nil { - break - } - - length := binary.BigEndian.Uint16(header[tls.SizeRecordType+tls.SizeVersion:]) - rec := make([]byte, length) - _, err := io.ReadFull(reader, rec) - suite.NoError(err) - - assembled.Write(rec) - } - - suite.Equal([]byte("firstsecond"), assembled.Bytes()) -} - -func (suite *ConnTestSuite) TestSyncWriteBlocksAsyncWrite() { - suite.connMock. - On("Write", mock.AnythingOfType("[]uint8")). - Return(0, nil). - Maybe() - - c := suite.makeConn() - defer c.Stop() - - // Start SyncWrite — it holds exclusive lock. - syncDone := make(chan struct{}) - - go func() { - defer close(syncDone) - c.SyncWrite([]byte("exclusive")) //nolint: errcheck - }() - - // Give SyncWrite time to acquire the lock. - time.Sleep(10 * time.Millisecond) - - // Async Write should block until SyncWrite completes. - writeDone := make(chan struct{}) - - go func() { - defer close(writeDone) - c.Write([]byte("blocked")) //nolint: errcheck - }() - - // SyncWrite should finish first. - <-syncDone - - select { - case <-writeDone: - // Write completed after SyncWrite — correct. - case <-time.After(2 * time.Second): - suite.Fail("async Write did not unblock after SyncWrite completed") - } -} - -func (suite *ConnTestSuite) TestSyncWriteReturnsErrorAfterStop() { - suite.connMock. - On("Write", mock.AnythingOfType("[]uint8")). - Return(0, nil). - Maybe() - - c := suite.makeConn() - c.Stop() - - time.Sleep(10 * time.Millisecond) - - _, err := c.SyncWrite([]byte("too late")) - suite.Error(err) -} - func TestConn(t *testing.T) { t.Parallel() suite.Run(t, &ConnTestSuite{}) From 724904f50d8310dfc26be6d31c0abba14137c306 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Thu, 19 Mar 2026 15:42:00 +0100 Subject: [PATCH 3/6] Wait in doppel.Conn if there is anything to write --- mtglib/internal/doppel/conn.go | 52 ++++++++++++++++++----------- mtglib/internal/doppel/conn_test.go | 31 +++++++++++++++++ 2 files changed, 63 insertions(+), 20 deletions(-) diff --git a/mtglib/internal/doppel/conn.go b/mtglib/internal/doppel/conn.go index bd76c7f..33ea88d 100644 --- a/mtglib/internal/doppel/conn.go +++ b/mtglib/internal/doppel/conn.go @@ -16,22 +16,25 @@ type Conn struct { } type connPayload struct { - ctx context.Context - ctxCancel context.CancelCauseFunc - clock Clock - wg sync.WaitGroup - syncWriteLock sync.RWMutex - writeStream bytes.Buffer - writeCond *sync.Cond + ctx context.Context + ctxCancel context.CancelCauseFunc + clock Clock + wg sync.WaitGroup + writeStream bytes.Buffer + writtenCond sync.Cond + done bool } func (c Conn) Write(p []byte) (int, error) { - c.p.syncWriteLock.RLock() - defer c.p.syncWriteLock.RUnlock() + if len(p) == 0 { + return 0, context.Cause(c.p.ctx) + } - c.p.writeCond.L.Lock() + c.p.writtenCond.L.Lock() c.p.writeStream.Write(p) - c.p.writeCond.L.Unlock() + c.p.writtenCond.L.Unlock() + + c.p.writtenCond.Signal() return len(p), context.Cause(c.p.ctx) } @@ -43,8 +46,6 @@ func (c Conn) Start() { } func (c Conn) start() { - defer c.p.writeCond.Broadcast() - buf := [tls.MaxRecordSize]byte{} for { @@ -54,11 +55,16 @@ func (c Conn) start() { case <-c.p.clock.tick: } - c.p.writeCond.L.Lock() - n, err := c.p.writeStream.Read(buf[:c.p.clock.stats.Size()]) - c.p.writeCond.L.Unlock() + size := c.p.clock.stats.Size() - if n == 0 || err != nil { + c.p.writtenCond.L.Lock() + for c.p.writeStream.Len() == 0 && !c.p.done { + c.p.writtenCond.Wait() + } + n, _ := c.p.writeStream.Read(buf[:size]) + c.p.writtenCond.L.Unlock() + + if n == 0 { continue } @@ -66,13 +72,17 @@ func (c Conn) start() { c.p.ctxCancel(err) return } - - c.p.writeCond.Signal() } } func (c Conn) Stop() { c.p.ctxCancel(nil) + + c.p.writtenCond.L.Lock() + c.p.done = true + c.p.writtenCond.L.Unlock() + c.p.writtenCond.Broadcast() + c.p.wg.Wait() } @@ -83,7 +93,9 @@ func NewConn(ctx context.Context, conn essentials.Conn, stats *Stats) Conn { p: &connPayload{ ctx: ctx, ctxCancel: cancel, - writeCond: sync.NewCond(&sync.Mutex{}), + writtenCond: sync.Cond{ + L: &sync.Mutex{}, + }, clock: Clock{ stats: stats, tick: make(chan struct{}), diff --git a/mtglib/internal/doppel/conn_test.go b/mtglib/internal/doppel/conn_test.go index e774bce..8501469 100644 --- a/mtglib/internal/doppel/conn_test.go +++ b/mtglib/internal/doppel/conn_test.go @@ -141,6 +141,37 @@ func (suite *ConnTestSuite) TestWriteReturnsErrorAfterStop() { suite.Error(err) } +func (suite *ConnTestSuite) TestStopDoesNotDeadlockWhenStartIsWaiting() { + suite.connMock. + On("Write", mock.AnythingOfType("[]uint8")). + Return(0, nil). + Maybe() + + for range 100 { + func() { + ctx, cancel := context.WithCancel(suite.ctx) + defer cancel() + + c := NewConn(ctx, suite.connMock, &Stats{ + k: 2.0, + lambda: 0.01, + }) + + done := make(chan struct{}) + go func() { + defer close(done) + c.Stop() + }() + + select { + case <-done: + case <-time.After(2 * time.Second): + suite.Fail("Stop() deadlocked: start() likely stuck in writtenCond.Wait()") + } + }() + } +} + func (suite *ConnTestSuite) TestStopOnUnderlyingWriteError() { suite.connMock. On("Write", mock.AnythingOfType("[]uint8")). From cb436efd87e959239d7da174d453b4f3ec7443a2 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Thu, 19 Mar 2026 17:37:51 +0100 Subject: [PATCH 4/6] Avoid double buffering in TLS hot path --- mtglib/internal/doppel/conn.go | 4 +- mtglib/internal/tls/utils.go | 24 ++++++---- mtglib/internal/tls/utils_test.go | 78 +++++++++++++++++++++++++++++++ 3 files changed, 94 insertions(+), 12 deletions(-) diff --git a/mtglib/internal/doppel/conn.go b/mtglib/internal/doppel/conn.go index 33ea88d..cc81d91 100644 --- a/mtglib/internal/doppel/conn.go +++ b/mtglib/internal/doppel/conn.go @@ -61,14 +61,14 @@ func (c Conn) start() { for c.p.writeStream.Len() == 0 && !c.p.done { c.p.writtenCond.Wait() } - n, _ := c.p.writeStream.Read(buf[:size]) + n, _ := c.p.writeStream.Read(buf[tls.SizeHeader : tls.SizeHeader+size]) c.p.writtenCond.L.Unlock() if n == 0 { continue } - if err := tls.WriteRecord(c.Conn, buf[:n]); err != nil { + if err := tls.WriteRecordInPlace(c.Conn, buf[:], n); err != nil { c.p.ctxCancel(err) return } diff --git a/mtglib/internal/tls/utils.go b/mtglib/internal/tls/utils.go index 978e048..c0bdc39 100644 --- a/mtglib/internal/tls/utils.go +++ b/mtglib/internal/tls/utils.go @@ -29,20 +29,24 @@ func ReadRecord(r io.Reader, w io.Writer) (byte, int64, error) { func WriteRecord(w io.Writer, payload []byte) error { buf := [MaxRecordSize]byte{} - buf[0] = TypeApplicationData + copy(buf[SizeHeader:], payload) - bufV := buf[SizeRecordType:] - copy(bufV[:SizeVersion], TLSVersion[:]) + return WriteRecordInPlace(w, buf[:], len(payload)) +} - bufS := bufV[SizeVersion:] - binary.BigEndian.PutUint16(bufS[:SizeSize], uint16(len(payload))) - - bufP := buf[SizeHeader:] - if n := copy(bufP, payload); n != len(payload) { - return fmt.Errorf("copied %d bytes of payload instead of %d", n, len(payload)) +func WriteRecordInPlace(w io.Writer, buf []byte, payloadLen int) error { + if payloadLen > MaxRecordPayloadSize { + return fmt.Errorf("payload %d exceeds max %d", payloadLen, MaxRecordPayloadSize) } - _, err := w.Write(buf[:SizeHeader+len(payload)]) + buf[0] = TypeApplicationData + copy(buf[SizeRecordType:SizeRecordType+SizeVersion], TLSVersion[:]) + binary.BigEndian.PutUint16( + buf[SizeRecordType+SizeVersion:SizeRecordType+SizeVersion+SizeSize], + uint16(payloadLen), + ) + + _, err := w.Write(buf[:SizeHeader+payloadLen]) return err } diff --git a/mtglib/internal/tls/utils_test.go b/mtglib/internal/tls/utils_test.go index 9ddbfa8..ad5a47f 100644 --- a/mtglib/internal/tls/utils_test.go +++ b/mtglib/internal/tls/utils_test.go @@ -119,6 +119,84 @@ func (suite *UtilsTestSuite) TestWriteRecordPayloadTooLarge() { suite.Error(err) } +func (suite *UtilsTestSuite) TestWriteRecordInPlace() { + payload := []byte("hello in-place") + + var buf [MaxRecordSize]byte + copy(buf[SizeHeader:], payload) + + err := WriteRecordInPlace(suite.dst, buf[:], len(payload)) + suite.NoError(err) + + written := suite.dst.Bytes() + suite.Equal(byte(TypeApplicationData), written[0]) + suite.Equal(TLSVersion[:], written[SizeRecordType:SizeRecordType+SizeVersion]) + + length := binary.BigEndian.Uint16(written[SizeRecordType+SizeVersion:]) + suite.Equal(uint16(len(payload)), length) + suite.Equal(payload, written[SizeHeader:]) +} + +func (suite *UtilsTestSuite) TestWriteRecordInPlaceRoundTrip() { + payload := []byte("round trip in-place") + + var buf [MaxRecordSize]byte + copy(buf[SizeHeader:], payload) + + var wire bytes.Buffer + + err := WriteRecordInPlace(&wire, buf[:], len(payload)) + suite.NoError(err) + + var recovered bytes.Buffer + + recordType, length, err := ReadRecord(&wire, &recovered) + suite.NoError(err) + suite.Equal(byte(TypeApplicationData), recordType) + suite.Equal(int64(len(payload)), length) + suite.Equal(payload, recovered.Bytes()) +} + +func (suite *UtilsTestSuite) TestWriteRecordInPlacePayloadTooLarge() { + var buf [MaxRecordSize]byte + + err := WriteRecordInPlace(suite.dst, buf[:], MaxRecordPayloadSize+1) + suite.Error(err) +} + +func (suite *UtilsTestSuite) TestWriteRecordInPlacePropagatesError() { + m := &WriterMock{} + m. + On("Write", mock.AnythingOfType("[]uint8")). + Once(). + Return(0, errors.New("disk full")) + + var buf [MaxRecordSize]byte + copy(buf[SizeHeader:], []byte("data")) + + err := WriteRecordInPlace(m, buf[:], 4) + suite.Error(err) + + m.AssertExpectations(suite.T()) +} + +func (suite *UtilsTestSuite) TestWriteRecordInPlaceMatchesWriteRecord() { + payload := []byte("equivalence check") + + var legacy bytes.Buffer + err := WriteRecord(&legacy, payload) + suite.NoError(err) + + var buf [MaxRecordSize]byte + copy(buf[SizeHeader:], payload) + + var inPlace bytes.Buffer + err = WriteRecordInPlace(&inPlace, buf[:], len(payload)) + suite.NoError(err) + + suite.Equal(legacy.Bytes(), inPlace.Bytes()) +} + func TestUtils(t *testing.T) { t.Parallel() suite.Run(t, &UtilsTestSuite{}) From feb57004e1acce3c948f29db4349f505c5c263b8 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Thu, 19 Mar 2026 17:39:48 +0100 Subject: [PATCH 5/6] Fix reslicing --- mtglib/internal/doppel/ganger.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/mtglib/internal/doppel/ganger.go b/mtglib/internal/doppel/ganger.go index 6621761..c8fbfbc 100644 --- a/mtglib/internal/doppel/ganger.go +++ b/mtglib/internal/doppel/ganger.go @@ -98,7 +98,8 @@ func (g *Ganger) run() { g.durations = append(g.durations, durations...) if len(g.durations) > DoppelGangerMaxDurations { - g.durations = g.durations[len(g.durations)-DoppelGangerMaxDurations:] + copy(g.durations, g.durations[len(g.durations)-DoppelGangerMaxDurations:]) + g.durations = g.durations[:DoppelGangerMaxDurations] } if len(g.durations) < MinDurationsToCalculate { From 4a8d099acadd9df48e580ea4c535400ab596e83d Mon Sep 17 00:00:00 2001 From: 9seconds Date: Thu, 19 Mar 2026 17:39:57 +0100 Subject: [PATCH 6/6] Remove unused buffer --- mtglib/internal/tls/conn.go | 2 -- 1 file changed, 2 deletions(-) diff --git a/mtglib/internal/tls/conn.go b/mtglib/internal/tls/conn.go index e889310..ce33fa1 100644 --- a/mtglib/internal/tls/conn.go +++ b/mtglib/internal/tls/conn.go @@ -34,7 +34,6 @@ type Conn struct { type connPayload struct { readBuf bytes.Buffer - writeBuf bytes.Buffer connBuffered *bufio.Reader read bool write bool @@ -80,7 +79,6 @@ func New(conn essentials.Conn, read, write bool) Conn { } newConn.p.readBuf.Grow(DefaultBufferSize) - newConn.p.writeBuf.Grow(DefaultBufferSize) return newConn }