aboutsummaryrefslogtreecommitdiff
path: root/hotline/client_conn_test.go
diff options
context:
space:
mode:
authorJeff Halter <868228+jhalter@users.noreply.github.com>2026-06-12 08:11:55 -0700
committerJeff Halter <868228+jhalter@users.noreply.github.com>2026-06-12 08:11:55 -0700
commit7ebc802d0a269218f05b3b51eb10ac66eacb4d1f (patch)
tree0c700edcde0fcc034605dfcf98b81ddb1cf174b6 /hotline/client_conn_test.go
parentd2791fcbadcf332dc05e5bebc7350eca50617263 (diff)
Replace shared outbox with per-client send queues
The outbox channel spawned one goroutine per outbound transaction, so concurrent sends to the same client could interleave bytes within the transaction framing, per-client message ordering was not guaranteed, and a slow client accumulated unbounded goroutines. Each ClientConn now has a bounded send queue drained by a single writer goroutine, which serializes writes and preserves enqueue order. Send never blocks: if a client's queue overflows, its connection is closed and the read loop performs the usual disconnect cleanup. Server.Send routes transactions to the target client's queue, replacing processOutbox and sendTransaction. Handler signatures are unchanged. Disconnect now removes the client from the manager before notifying peers so no new transactions are routed to a departing client, then idempotently closes its send queue. New tests cover write ordering, framing integrity under concurrent senders, the slow-client disconnect policy, and a Send/Disconnect race exercise (run with -race).
Diffstat (limited to 'hotline/client_conn_test.go')
-rw-r--r--hotline/client_conn_test.go220
1 files changed, 188 insertions, 32 deletions
diff --git a/hotline/client_conn_test.go b/hotline/client_conn_test.go
index 94e38e9..5c5463d 100644
--- a/hotline/client_conn_test.go
+++ b/hotline/client_conn_test.go
@@ -1,7 +1,10 @@
package hotline
import (
+ "bufio"
"bytes"
+ "fmt"
+ "sync"
"testing"
"github.com/stretchr/testify/assert"
@@ -343,13 +346,11 @@ func TestClientConn_Disconnect(t *testing.T) {
mockMgr.On("Delete", ClientID{0, 1}).Return()
mockMgr.On("List").Return([]*ClientConn{})
- outbox := make(chan Transaction, 10)
cc := &ClientConn{
ID: ClientID{0, 1},
Connection: &nopCloserRWC{Buffer: &bytes.Buffer{}},
Server: &Server{
ClientMgr: mockMgr,
- outbox: outbox,
Logger: NewTestLogger(),
},
}
@@ -357,39 +358,41 @@ func TestClientConn_Disconnect(t *testing.T) {
cc.Disconnect()
mockMgr.AssertCalled(t, "Delete", ClientID{0, 1})
- assert.Empty(t, outbox) // No other clients to notify
})
t.Run("notifies other clients", func(t *testing.T) {
+ peer2 := &ClientConn{ID: ClientID{0, 2}}
+ peer3 := &ClientConn{ID: ClientID{0, 3}}
+
mockMgr := &MockClientMgr{}
mockMgr.On("Delete", ClientID{0, 1}).Return()
mockMgr.On("List").Return([]*ClientConn{
{ID: ClientID{0, 1}},
- {ID: ClientID{0, 2}},
- {ID: ClientID{0, 3}},
+ peer2,
+ peer3,
})
+ mockMgr.On("Get", ClientID{0, 2}).Return(peer2)
+ mockMgr.On("Get", ClientID{0, 3}).Return(peer3)
- outbox := make(chan Transaction, 10)
cc := &ClientConn{
ID: ClientID{0, 1},
Connection: &nopCloserRWC{Buffer: &bytes.Buffer{}},
Server: &Server{
ClientMgr: mockMgr,
- outbox: outbox,
Logger: NewTestLogger(),
},
}
cc.Disconnect()
- assert.Len(t, outbox, 2)
+ assert.Len(t, peer2.sendCh, 1)
+ assert.Len(t, peer3.sendCh, 1)
mockMgr.AssertExpectations(t)
})
}
func TestClientConn_handleTransaction(t *testing.T) {
t.Run("dispatches to registered handler", func(t *testing.T) {
- outbox := make(chan Transaction, 10)
mockMgr := &MockClientMgr{}
cc := &ClientConn{
@@ -397,7 +400,6 @@ func TestClientConn_handleTransaction(t *testing.T) {
Account: &Account{},
Logger: NewTestLogger(),
Server: &Server{
- outbox: outbox,
ClientMgr: mockMgr,
handlers: map[TranType]HandlerFunc{
TranChatSend: func(cc *ClientConn, t *Transaction) []Transaction {
@@ -406,15 +408,16 @@ func TestClientConn_handleTransaction(t *testing.T) {
},
},
}
+ mockMgr.On("Get", ClientID{0, 1}).Return(cc)
cc.handleTransaction(NewTransaction(TranChatSend, ClientID{0, 1}))
- assert.Len(t, outbox, 1)
+ assert.Len(t, cc.sendCh, 1)
assert.Equal(t, 0, cc.IdleTime)
})
t.Run("keepalive does not reset idle time", func(t *testing.T) {
- outbox := make(chan Transaction, 10)
+ mockMgr := &MockClientMgr{}
cc := &ClientConn{
ID: ClientID{0, 1},
@@ -422,7 +425,7 @@ func TestClientConn_handleTransaction(t *testing.T) {
IdleTime: 100,
Logger: NewTestLogger(),
Server: &Server{
- outbox: outbox,
+ ClientMgr: mockMgr,
handlers: map[TranType]HandlerFunc{
TranKeepAlive: func(cc *ClientConn, t *Transaction) []Transaction {
return []Transaction{cc.NewReply(t)}
@@ -430,6 +433,7 @@ func TestClientConn_handleTransaction(t *testing.T) {
},
},
}
+ mockMgr.On("Get", ClientID{0, 1}).Return(cc)
cc.handleTransaction(NewTransaction(TranKeepAlive, ClientID{0, 1}))
@@ -437,11 +441,9 @@ func TestClientConn_handleTransaction(t *testing.T) {
})
t.Run("non-keepalive clears away flag", func(t *testing.T) {
- outbox := make(chan Transaction, 10)
+ peer := &ClientConn{ID: ClientID{0, 1}}
mockMgr := &MockClientMgr{}
- mockMgr.On("List").Return([]*ClientConn{
- {ID: ClientID{0, 1}},
- })
+ mockMgr.On("List").Return([]*ClientConn{peer})
cc := &ClientConn{
ID: ClientID{0, 1},
@@ -451,7 +453,6 @@ func TestClientConn_handleTransaction(t *testing.T) {
IdleTime: 50,
Logger: NewTestLogger(),
Server: &Server{
- outbox: outbox,
ClientMgr: mockMgr,
handlers: map[TranType]HandlerFunc{
TranChatSend: func(cc *ClientConn, t *Transaction) []Transaction {
@@ -467,40 +468,195 @@ func TestClientConn_handleTransaction(t *testing.T) {
assert.Equal(t, 0, cc.IdleTime)
assert.False(t, cc.Flags.IsSet(UserFlagAway))
// SendAll should have sent TranNotifyChangeUser
- assert.Greater(t, len(outbox), 0)
+ assert.Greater(t, len(peer.sendCh), 0)
})
}
func TestClientConn_SendAll(t *testing.T) {
- mockMgr := &MockClientMgr{}
- mockMgr.On("List").Return([]*ClientConn{
+ peers := []*ClientConn{
{ID: ClientID{0, 1}},
{ID: ClientID{0, 2}},
{ID: ClientID{0, 3}},
- })
+ }
+ mockMgr := &MockClientMgr{}
+ mockMgr.On("List").Return(peers)
- outbox := make(chan Transaction, 10)
cc := &ClientConn{
ID: ClientID{0, 1},
Server: &Server{
ClientMgr: mockMgr,
- outbox: outbox,
},
}
cc.SendAll(TranChatMsg, NewField(FieldData, []byte("hello")))
- assert.Len(t, outbox, 3)
+ for _, peer := range peers {
+ assert.Len(t, peer.sendCh, 1)
- clientIDs := make(map[ClientID]bool)
- for range 3 {
- tran := <-outbox
- clientIDs[tran.ClientID] = true
+ tran := <-peer.sendCh
+ assert.Equal(t, peer.ID, tran.ClientID)
assert.Equal(t, TranChatMsg, tran.Type)
}
- assert.True(t, clientIDs[ClientID{0, 1}])
- assert.True(t, clientIDs[ClientID{0, 2}])
- assert.True(t, clientIDs[ClientID{0, 3}])
mockMgr.AssertExpectations(t)
}
+
+// closeRecorderRWC records whether Close was called.
+type closeRecorderRWC struct {
+ *bytes.Buffer
+ closed bool
+}
+
+func (c *closeRecorderRWC) Close() error {
+ c.closed = true
+ return nil
+}
+
+// TestClientConn_writeLoop_ordering verifies that transactions are written to the connection in
+// the order they were enqueued.
+func TestClientConn_writeLoop_ordering(t *testing.T) {
+ buf := &bytes.Buffer{}
+ cc := &ClientConn{
+ ID: ClientID{0, 1},
+ Connection: &nopCloserRWC{Buffer: buf},
+ Logger: NewTestLogger(),
+ }
+
+ done := make(chan struct{})
+ go func() {
+ defer close(done)
+ cc.writeLoop()
+ }()
+
+ const numTrans = 50
+ for i := range numTrans {
+ cc.Send(NewTransaction(TranChatMsg, cc.ID, NewField(FieldData, fmt.Appendf(nil, "msg-%03d", i))))
+ }
+
+ // Closing the queue stops writeLoop after it drains the remaining transactions.
+ cc.closeSendQueue()
+ <-done
+
+ scanner := bufio.NewScanner(bytes.NewReader(buf.Bytes()))
+ scanner.Split(transactionScanner)
+
+ var count int
+ for scanner.Scan() {
+ var tran Transaction
+ _, err := tran.Write(scanner.Bytes())
+ require.NoError(t, err, "transaction %d is malformed", count)
+
+ assert.Equal(t, fmt.Sprintf("msg-%03d", count), string(tran.GetField(FieldData).Data))
+ count++
+ }
+ require.NoError(t, scanner.Err())
+ assert.Equal(t, numTrans, count)
+}
+
+// TestClientConn_writeLoop_noInterleaving verifies that concurrent senders cannot interleave bytes
+// within the connection's transaction framing.
+func TestClientConn_writeLoop_noInterleaving(t *testing.T) {
+ buf := &bytes.Buffer{}
+ cc := &ClientConn{
+ ID: ClientID{0, 1},
+ Connection: &nopCloserRWC{Buffer: buf},
+ Logger: NewTestLogger(),
+ }
+
+ done := make(chan struct{})
+ go func() {
+ defer close(done)
+ cc.writeLoop()
+ }()
+
+ // Total sends must stay within sendQueueDepth so the queue cannot overflow even if writeLoop
+ // has not started draining yet.
+ const senders, transPerSender = 4, 10
+
+ var wg sync.WaitGroup
+ for s := range senders {
+ wg.Add(1)
+ go func() {
+ defer wg.Done()
+ for i := range transPerSender {
+ cc.Send(NewTransaction(TranChatMsg, cc.ID, NewField(FieldData, fmt.Appendf(nil, "sender-%d-msg-%d", s, i))))
+ }
+ }()
+ }
+ wg.Wait()
+
+ cc.closeSendQueue()
+ <-done
+
+ scanner := bufio.NewScanner(bytes.NewReader(buf.Bytes()))
+ scanner.Split(transactionScanner)
+
+ var count int
+ for scanner.Scan() {
+ var tran Transaction
+ _, err := tran.Write(scanner.Bytes())
+ require.NoError(t, err, "transaction %d is malformed", count)
+ assert.Equal(t, TranChatMsg, tran.Type)
+ count++
+ }
+ require.NoError(t, scanner.Err())
+ assert.Equal(t, senders*transPerSender, count)
+}
+
+// TestClientConn_Send_slowClientDisconnected verifies that a client whose send queue overflows has
+// its connection closed and that subsequent sends are dropped without panicking.
+func TestClientConn_Send_slowClientDisconnected(t *testing.T) {
+ conn := &closeRecorderRWC{Buffer: &bytes.Buffer{}}
+ cc := &ClientConn{
+ ID: ClientID{0, 1},
+ Connection: conn,
+ Logger: NewTestLogger(),
+ }
+
+ // No writeLoop is running, so the queue fills up after sendQueueDepth sends.
+ for i := range sendQueueDepth + 1 {
+ cc.Send(NewTransaction(TranChatMsg, cc.ID, NewField(FieldData, fmt.Appendf(nil, "msg-%d", i))))
+ }
+
+ assert.True(t, conn.closed, "connection should be closed when the send queue overflows")
+
+ // Further sends must be silently dropped.
+ cc.Send(NewTransaction(TranChatMsg, cc.ID))
+}
+
+// TestClientConn_SendDisconnectRace exercises concurrent Send and Disconnect calls. Run with
+// -race to detect data races and unsynchronized channel closes.
+func TestClientConn_SendDisconnectRace(t *testing.T) {
+ mockMgr := &MockClientMgr{}
+ mockMgr.On("Delete", ClientID{0, 1}).Return()
+ mockMgr.On("List").Return([]*ClientConn{})
+
+ cc := &ClientConn{
+ ID: ClientID{0, 1},
+ Connection: &nopCloserRWC{Buffer: &bytes.Buffer{}},
+ Logger: NewTestLogger(),
+ Server: &Server{
+ ClientMgr: mockMgr,
+ Logger: NewTestLogger(),
+ },
+ }
+
+ var wg sync.WaitGroup
+ for range 4 {
+ wg.Add(1)
+ go func() {
+ defer wg.Done()
+ for i := range 100 {
+ cc.Send(NewTransaction(TranChatMsg, cc.ID, NewField(FieldData, fmt.Appendf(nil, "msg-%d", i))))
+ }
+ }()
+ }
+
+ wg.Add(1)
+ go func() {
+ defer wg.Done()
+ cc.Disconnect()
+ }()
+
+ wg.Wait()
+}