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) } }