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)