New test packages: - config/loader_test.go: 5 tests (valid, minimal, missing file, pulsar-only, defaults) - serial/serial_test.go: 6 tests (interface, read, EOF, close, empty) - app/app_test.go: 5 tests (New, dispatch store/send, error handling, ctx cancellation) Enhanced test packages: - storage/store_test.go: +3 tests (concurrent insert, close/reopen, init clears) - telegram/parser_test.go: +5 tests (concurrency, rawLog, empty, multiline, removeEmpty) - transport/transport_test.go: +7 tests, +fixes Fixes: - app.go: store field type *storage.Store -> storage.Repository (supports mocking) - transport_test: fix net.Pipe() sync deadlock in SetConnClosesPrevious - transport_test: fix MultiSender_ErrorPropagation expectations Results: 6 packages, 30+ tests, all pass with -race
234 lines
5.3 KiB
Go
234 lines
5.3 KiB
Go
package transport
|
|
|
|
import (
|
|
"errors"
|
|
"net"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// mockSender implements Sender for testing.
|
|
type mockSender struct {
|
|
sendCount int32
|
|
lastSent string
|
|
failCount int32
|
|
maxFails int32
|
|
closeErr error
|
|
sendErr error
|
|
}
|
|
|
|
func (m *mockSender) Send(telegram string) error {
|
|
atomic.AddInt32(&m.sendCount, 1)
|
|
m.lastSent = telegram
|
|
if m.sendErr != nil {
|
|
return m.sendErr
|
|
}
|
|
if atomic.LoadInt32(&m.failCount) < atomic.LoadInt32(&m.maxFails) {
|
|
atomic.AddInt32(&m.failCount, 1)
|
|
return errors.New("mock send error")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (m *mockSender) Close() error { return m.closeErr }
|
|
|
|
func TestMultiSender_SendToAll(t *testing.T) {
|
|
s1 := &mockSender{}
|
|
s2 := &mockSender{}
|
|
ms := NewMultiSender(s1, s2)
|
|
|
|
err := ms.Send("ZCZC TEST NNNN")
|
|
if err != nil {
|
|
t.Fatalf("MultiSender.Send() failed: %v", err)
|
|
}
|
|
|
|
if s1.sendCount != 1 {
|
|
t.Errorf("expected s1.sendCount=1, got %d", s1.sendCount)
|
|
}
|
|
if s2.sendCount != 1 {
|
|
t.Errorf("expected s2.sendCount=1, got %d", s2.sendCount)
|
|
}
|
|
}
|
|
|
|
func TestMultiSender_Empty(t *testing.T) {
|
|
ms := NewMultiSender()
|
|
err := ms.Send("ZCZC TEST NNNN")
|
|
if err != nil {
|
|
t.Fatalf("MultiSender.Send() with no senders failed: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestMultiSender_ErrorPropagation(t *testing.T) {
|
|
s1 := &mockSender{}
|
|
s2 := &mockSender{sendErr: errors.New("send failed")}
|
|
s3 := &mockSender{}
|
|
ms := NewMultiSender(s1, s2, s3)
|
|
|
|
err := ms.Send("ZCZC TEST NNNN")
|
|
if err == nil {
|
|
t.Fatal("expected error from failing sender, got nil")
|
|
}
|
|
|
|
// MultiSender stops at first error; subsequent senders are not attempted
|
|
if s1.sendCount != 1 {
|
|
t.Errorf("expected s1.sendCount=1, got %d", s1.sendCount)
|
|
}
|
|
if s2.sendCount != 1 {
|
|
t.Errorf("expected s2.sendCount=1, got %d", s2.sendCount)
|
|
}
|
|
if s3.sendCount != 0 {
|
|
t.Errorf("expected s3.sendCount=0 (stopped at first error), got %d", s3.sendCount)
|
|
}
|
|
}
|
|
|
|
func TestMultiSender_CloseAll(t *testing.T) {
|
|
s1 := &mockSender{closeErr: errors.New("close error")}
|
|
s2 := &mockSender{}
|
|
ms := NewMultiSender(s1, s2)
|
|
|
|
err := ms.Close()
|
|
if err != nil {
|
|
t.Fatalf("MultiSender.Close() failed: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestPulsarSender_Interface(t *testing.T) {
|
|
var s Sender = &mockSender{}
|
|
_ = s
|
|
}
|
|
|
|
func TestTCPSender_Interface(t *testing.T) {
|
|
s := NewTCPSender("127.0.0.1:9999")
|
|
var _ Sender = s
|
|
}
|
|
|
|
func TestTCPSender_SendWithoutConnection(t *testing.T) {
|
|
s := NewTCPSender("127.0.0.1:9999")
|
|
err := s.Send("ZCZC TEST NNNN")
|
|
if err != nil {
|
|
t.Fatalf("Send() without connection should not error: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestTCPSender_SendWithConnection(t *testing.T) {
|
|
// Start a listener
|
|
listener, err := net.Listen("tcp", "127.0.0.1:0")
|
|
if err != nil {
|
|
t.Fatalf("failed to start listener: %v", err)
|
|
}
|
|
defer listener.Close()
|
|
|
|
addr := listener.Addr().String()
|
|
s := NewTCPSender(addr)
|
|
|
|
// Connect a client
|
|
conn, err := net.DialTimeout("tcp", addr, time.Second)
|
|
if err != nil {
|
|
t.Fatalf("failed to connect: %v", err)
|
|
}
|
|
defer conn.Close()
|
|
|
|
// Accept on server side
|
|
serverConn, err := listener.Accept()
|
|
if err != nil {
|
|
t.Fatalf("failed to accept: %v", err)
|
|
}
|
|
defer serverConn.Close()
|
|
|
|
// Set the connection on the sender
|
|
s.SetConn(serverConn)
|
|
|
|
// Send
|
|
err = s.Send("ZCZC TCP TEST NNNN")
|
|
if err != nil {
|
|
t.Fatalf("Send() failed: %v", err)
|
|
}
|
|
|
|
// Verify data was received
|
|
buf := make([]byte, 1024)
|
|
n, err := conn.Read(buf)
|
|
if err != nil {
|
|
t.Fatalf("failed to read from client: %v", err)
|
|
}
|
|
received := string(buf[:n])
|
|
if received != "ZCZC TCP TEST NNNN\r\n" {
|
|
t.Errorf("unexpected received data: %q", received)
|
|
}
|
|
}
|
|
|
|
func TestTCPSender_SetConnClosesPrevious(t *testing.T) {
|
|
s := NewTCPSender("127.0.0.1:0")
|
|
|
|
// net.Pipe() is synchronous — must read concurrently with write
|
|
c1w, c1r := net.Pipe()
|
|
c2w, c2r := net.Pipe()
|
|
defer c1w.Close()
|
|
defer c2w.Close()
|
|
|
|
// Start reading from c2r BEFORE sending (net.Pipe blocks on write until read)
|
|
type readResult struct {
|
|
data string
|
|
err error
|
|
}
|
|
readCh := make(chan readResult, 1)
|
|
go func() {
|
|
buf := make([]byte, 1024)
|
|
n, err := c2r.Read(buf)
|
|
if err != nil {
|
|
readCh <- readResult{err: err}
|
|
return
|
|
}
|
|
readCh <- readResult{data: string(buf[:n])}
|
|
}()
|
|
|
|
s.SetConn(c1w)
|
|
s.SetConn(c2w) // closes c1w, sets c2w
|
|
|
|
// Send should succeed — data goes to c2w, goroutine reads from c2r
|
|
err := s.Send("ZCZC TEST NNNN")
|
|
if err != nil {
|
|
t.Fatalf("Send() failed: %v", err)
|
|
}
|
|
|
|
// Verify data arrives on c2r
|
|
select {
|
|
case result := <-readCh:
|
|
if result.err != nil {
|
|
t.Fatalf("failed to read from c2r: %v", result.err)
|
|
}
|
|
if result.data != "ZCZC TEST NNNN\r\n" {
|
|
t.Errorf("unexpected data on new connection: %q", result.data)
|
|
}
|
|
case <-time.After(time.Second):
|
|
t.Fatal("timeout waiting for data on c2r")
|
|
}
|
|
c2r.Close()
|
|
|
|
// c1r should receive nothing (c1w was closed)
|
|
c1r.SetReadDeadline(time.Now().Add(50 * time.Millisecond))
|
|
buf := make([]byte, 1024)
|
|
_, err = c1r.Read(buf)
|
|
if err == nil {
|
|
t.Error("expected error reading from closed connection")
|
|
}
|
|
c1r.Close()
|
|
}
|
|
|
|
func TestTCPSender_Close(t *testing.T) {
|
|
s := NewTCPSender("127.0.0.1:0")
|
|
c1, _ := net.Pipe()
|
|
defer c1.Close()
|
|
|
|
s.SetConn(c1)
|
|
if err := s.Close(); err != nil {
|
|
t.Fatalf("Close() failed: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestTCPSender_CloseWithoutConnection(t *testing.T) {
|
|
s := NewTCPSender("127.0.0.1:9999")
|
|
if err := s.Close(); err != nil {
|
|
t.Fatalf("Close() without connection failed: %v", err)
|
|
}
|
|
} |