277 lines
6.1 KiB
Go
277 lines
6.1 KiB
Go
package transport
|
|||
|
|
|
||
|
|
import (
|
||
|
|
"context"
|
||
|
|
"io"
|
||
|
|
"net"
|
||
|
|
"sync"
|
||
|
|
"time"
|
||
|
|
|
||
|
|
"github.com/apache/pulsar-client-go/pulsar"
|
||
|
|
)
|
||
|
|
|
||
|
|
// Sender defines the interface for sending telegrams.
|
||
|
|
type Sender interface {
|
||
|
|
Send(telegram string) error
|
||
|
|
Close() error
|
||
|
|
}
|
||
|
|
|
||
|
|
// ─── TCP Sender ──────────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
// TCPSender sends telegrams over a TCP connection.
|
||
|
|
type TCPSender struct {
|
||
|
|
mu sync.RWMutex
|
||
|
|
conn net.Conn
|
||
|
|
addr string
|
||
|
|
}
|
||
|
|
|
||
|
|
// NewTCPSender creates a TCPSender. It does not connect immediately;
|
||
|
|
// calling Send will connect on demand.
|
||
|
|
func NewTCPSender(addr string) *TCPSender {
|
||
|
|
return &TCPSender{addr: addr}
|
||
|
|
}
|
||
|
|
|
||
|
|
// Send writes a telegram to the TCP connection.
|
||
|
|
func (s *TCPSender) Send(telegram string) error {
|
||
|
|
s.mu.RLock()
|
||
|
|
conn := s.conn
|
||
|
|
s.mu.RUnlock()
|
||
|
|
|
||
|
|
if conn == nil {
|
||
|
|
return nil // no client connected, silent drop
|
||
|
|
}
|
||
|
|
|
||
|
|
_, err := conn.Write([]byte(telegram + "\r\n"))
|
||
|
|
if err != nil {
|
||
|
|
s.mu.Lock()
|
||
|
|
s.conn.Close()
|
||
|
|
s.conn = nil
|
||
|
|
s.mu.Unlock()
|
||
|
|
}
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
// SetConn updates the active TCP connection.
|
||
|
|
func (s *TCPSender) SetConn(conn net.Conn) {
|
||
|
|
s.mu.Lock()
|
||
|
|
if s.conn != nil {
|
||
|
|
s.conn.Close()
|
||
|
|
}
|
||
|
|
s.conn = conn
|
||
|
|
s.mu.Unlock()
|
||
|
|
}
|
||
|
|
|
||
|
|
// Close closes the TCP connection.
|
||
|
|
func (s *TCPSender) Close() error {
|
||
|
|
s.mu.Lock()
|
||
|
|
defer s.mu.Unlock()
|
||
|
|
if s.conn != nil {
|
||
|
|
return s.conn.Close()
|
||
|
|
}
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// ─── TCP Server ──────────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
// TCPServer listens for TCP client connections. When a client connects,
|
||
|
|
// it updates the TCPSender's connection and replays unprocessed telegrams.
|
||
|
|
type TCPServer struct {
|
||
|
|
addr string
|
||
|
|
sender *TCPSender
|
||
|
|
store UnprocessedLoader
|
||
|
|
listener net.Listener
|
||
|
|
wg sync.WaitGroup
|
||
|
|
}
|
||
|
|
|
||
|
|
// TelegramRecord represents a stored telegram for replay.
|
||
|
|
type TelegramRecord struct {
|
||
|
|
ID int64
|
||
|
|
Text string
|
||
|
|
}
|
||
|
|
|
||
|
|
// UnprocessedLoader allows the TCP server to replay stored telegrams.
|
||
|
|
type UnprocessedLoader interface {
|
||
|
|
LoadUnprocessed() ([]TelegramRecord, error)
|
||
|
|
MarkProcessed(id int64) error
|
||
|
|
}
|
||
|
|
|
||
|
|
// NewTCPServer creates a TCP server that feeds connections into a TCPSender.
|
||
|
|
func NewTCPServer(addr string, sender *TCPSender, store UnprocessedLoader) *TCPServer {
|
||
|
|
return &TCPServer{
|
||
|
|
addr: addr,
|
||
|
|
sender: sender,
|
||
|
|
store: store,
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// Run starts the TCP listener loop. Blocks until ctx is cancelled.
|
||
|
|
func (s *TCPServer) Run(ctx context.Context) error {
|
||
|
|
var err error
|
||
|
|
s.listener, err = net.Listen("tcp", s.addr)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
failCount := 0
|
||
|
|
const maxFail = 3
|
||
|
|
|
||
|
|
for {
|
||
|
|
conn, err := s.listener.Accept()
|
||
|
|
if err != nil {
|
||
|
|
failCount++
|
||
|
|
if failCount >= maxFail {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
select {
|
||
|
|
case <-ctx.Done():
|
||
|
|
return ctx.Err()
|
||
|
|
case <-time.After(time.Second):
|
||
|
|
}
|
||
|
|
continue
|
||
|
|
}
|
||
|
|
failCount = 0
|
||
|
|
|
||
|
|
s.sender.SetConn(conn)
|
||
|
|
s.replayUnprocessed()
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// Stop shuts down the TCP listener.
|
||
|
|
func (s *TCPServer) Stop() {
|
||
|
|
if s.listener != nil {
|
||
|
|
s.listener.Close()
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *TCPServer) replayUnprocessed() {
|
||
|
|
if s.store == nil {
|
||
|
|
return
|
||
|
|
}
|
||
|
|
// Simplified: in production, load and send unprocessed telegrams
|
||
|
|
}
|
||
|
|
|
||
|
|
// ─── Pulsar Sender ───────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
// PulsarSender sends telegrams to Apache Pulsar.
|
||
|
|
type PulsarSender struct {
|
||
|
|
mu sync.Mutex
|
||
|
|
client pulsar.Client
|
||
|
|
producer pulsar.Producer
|
||
|
|
url string
|
||
|
|
topic string
|
||
|
|
name string
|
||
|
|
}
|
||
|
|
|
||
|
|
// NewPulsarSender creates a PulsarSender and initializes the producer.
|
||
|
|
func NewPulsarSender(url, topic, name string) (*PulsarSender, error) {
|
||
|
|
ps := &PulsarSender{url: url, topic: topic, name: name}
|
||
|
|
if err := ps.connect(); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
return ps, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (ps *PulsarSender) connect() error {
|
||
|
|
client, err := pulsar.NewClient(pulsar.ClientOptions{URL: ps.url})
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
producer, err := client.CreateProducer(pulsar.ProducerOptions{
|
||
|
|
Topic: ps.topic,
|
||
|
|
Name: ps.name,
|
||
|
|
})
|
||
|
|
if err != nil {
|
||
|
|
client.Close()
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
ps.mu.Lock()
|
||
|
|
ps.client = client
|
||
|
|
ps.producer = producer
|
||
|
|
ps.mu.Unlock()
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// Send publishes a telegram to Pulsar.
|
||
|
|
func (ps *PulsarSender) Send(telegram string) error {
|
||
|
|
ps.mu.Lock()
|
||
|
|
producer := ps.producer
|
||
|
|
ps.mu.Unlock()
|
||
|
|
|
||
|
|
var err error
|
||
|
|
for i := 0; i < 5; i++ {
|
||
|
|
_, err = producer.Send(context.Background(), &pulsar.ProducerMessage{
|
||
|
|
Payload: []byte(telegram),
|
||
|
|
})
|
||
|
|
if err == nil {
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
ps.reconnect()
|
||
|
|
}
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
func (ps *PulsarSender) reconnect() {
|
||
|
|
ps.mu.Lock()
|
||
|
|
if ps.producer != nil {
|
||
|
|
ps.producer.Close()
|
||
|
|
}
|
||
|
|
if ps.client != nil {
|
||
|
|
ps.client.Close()
|
||
|
|
}
|
||
|
|
ps.mu.Unlock()
|
||
|
|
ps.connect()
|
||
|
|
}
|
||
|
|
|
||
|
|
// Close closes the Pulsar producer and client.
|
||
|
|
func (ps *PulsarSender) Close() error {
|
||
|
|
ps.mu.Lock()
|
||
|
|
defer ps.mu.Unlock()
|
||
|
|
|
||
|
|
if ps.producer != nil {
|
||
|
|
ps.producer.Close()
|
||
|
|
}
|
||
|
|
if ps.client != nil {
|
||
|
|
ps.client.Close()
|
||
|
|
}
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// ─── MultiSender ─────────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
// MultiSender fans out telegrams to multiple Sender implementations.
|
||
|
|
type MultiSender struct {
|
||
|
|
senders []Sender
|
||
|
|
}
|
||
|
|
|
||
|
|
// NewMultiSender creates a MultiSender.
|
||
|
|
func NewMultiSender(senders ...Sender) *MultiSender {
|
||
|
|
return &MultiSender{senders: senders}
|
||
|
|
}
|
||
|
|
|
||
|
|
// Send sends to all registered senders. Errors are collected but all
|
||
|
|
// senders are attempted.
|
||
|
|
func (m *MultiSender) Send(telegram string) error {
|
||
|
|
for _, s := range m.senders {
|
||
|
|
if err := s.Send(telegram); err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// Close closes all registered senders.
|
||
|
|
func (m *MultiSender) Close() error {
|
||
|
|
for _, s := range m.senders {
|
||
|
|
s.Close()
|
||
|
|
}
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// Ensure interfaces are satisfied.
|
||
|
|
var _ Sender = (*TCPSender)(nil)
|
||
|
|
var _ Sender = (*PulsarSender)(nil)
|
||
|
|
var _ Sender = (*MultiSender)(nil)
|
||
|
|
var _ io.Closer = (*TCPSender)(nil)
|
||
|
|
var _ io.Closer = (*PulsarSender)(nil)
|