server.go 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302
  1. package icq_legacy
  2. import (
  3. "context"
  4. "errors"
  5. "fmt"
  6. "log/slog"
  7. "net"
  8. "sync"
  9. "time"
  10. "github.com/mk6i/open-oscar-server/config"
  11. )
  12. const (
  13. // MaxPacketSize is the maximum UDP packet size we'll accept
  14. MaxPacketSize = 8192
  15. // ReadBufferSize is the UDP socket read buffer size
  16. ReadBufferSize = 65536
  17. )
  18. var (
  19. ErrServerClosed = errors.New("server closed")
  20. ErrUnsupportedProto = errors.New("unsupported protocol version")
  21. )
  22. // LegacyServer handles legacy ICQ protocol connections over UDP
  23. type LegacyServer struct {
  24. conn *net.UDPConn
  25. config config.ICQLegacyConfig
  26. sessions *LegacySessionManager
  27. dispatcher *ProtocolDispatcher
  28. logger *slog.Logger
  29. stopChan chan struct{}
  30. wg sync.WaitGroup
  31. mu sync.RWMutex
  32. running bool
  33. }
  34. // NewLegacyServer creates a new legacy ICQ server
  35. func NewLegacyServer(
  36. cfg config.ICQLegacyConfig,
  37. sessions *LegacySessionManager,
  38. dispatcher *ProtocolDispatcher,
  39. logger *slog.Logger,
  40. ) *LegacyServer {
  41. return &LegacyServer{
  42. config: cfg,
  43. sessions: sessions,
  44. dispatcher: dispatcher,
  45. logger: logger,
  46. stopChan: make(chan struct{}),
  47. }
  48. }
  49. // Start begins listening for legacy ICQ connections
  50. func (s *LegacyServer) Start(ctx context.Context) error {
  51. if !s.config.Enabled {
  52. s.logger.Info("legacy ICQ server disabled")
  53. return nil
  54. }
  55. s.mu.Lock()
  56. if s.running {
  57. s.mu.Unlock()
  58. return errors.New("server already running")
  59. }
  60. s.running = true
  61. s.mu.Unlock()
  62. // Parse listen address
  63. addr, err := net.ResolveUDPAddr("udp4", s.config.UDPListener)
  64. if err != nil {
  65. return fmt.Errorf("invalid UDP listener address: %w", err)
  66. }
  67. // Create UDP socket
  68. conn, err := net.ListenUDP("udp4", addr)
  69. if err != nil {
  70. return fmt.Errorf("failed to listen on UDP: %w", err)
  71. }
  72. // Set socket options
  73. if err := conn.SetReadBuffer(ReadBufferSize); err != nil {
  74. s.logger.Warn("failed to set read buffer size", "err", err)
  75. }
  76. s.conn = conn
  77. s.logger.Info("legacy ICQ server started",
  78. "address", s.config.UDPListener,
  79. "versions", s.config.SupportedVersions,
  80. )
  81. // Start session cleanup routine
  82. s.wg.Add(1)
  83. go func() {
  84. defer s.wg.Done()
  85. s.sessions.StartCleanupRoutine(s.config.SessionTimeout/2, s.stopChan)
  86. }()
  87. // Start packet receive loop
  88. s.wg.Add(1)
  89. go func() {
  90. defer s.wg.Done()
  91. s.receiveLoop(ctx)
  92. }()
  93. return nil
  94. }
  95. // Stop gracefully shuts down the server
  96. func (s *LegacyServer) Stop() error {
  97. s.mu.Lock()
  98. if !s.running {
  99. s.mu.Unlock()
  100. return nil
  101. }
  102. s.running = false
  103. s.mu.Unlock()
  104. // Signal all goroutines to stop
  105. close(s.stopChan)
  106. // Close the UDP socket to unblock the receive loop
  107. if s.conn != nil {
  108. _ = s.conn.Close()
  109. }
  110. // Wait for all goroutines to finish
  111. s.wg.Wait()
  112. s.logger.Info("legacy ICQ server stopped")
  113. return nil
  114. }
  115. // receiveLoop handles incoming UDP packets
  116. func (s *LegacyServer) receiveLoop(ctx context.Context) {
  117. buf := make([]byte, MaxPacketSize)
  118. for {
  119. select {
  120. case <-ctx.Done():
  121. return
  122. case <-s.stopChan:
  123. return
  124. default:
  125. }
  126. // Set read deadline to allow periodic checking of stop signal
  127. _ = s.conn.SetReadDeadline(time.Now().Add(1 * time.Second))
  128. n, addr, err := s.conn.ReadFromUDP(buf)
  129. if err != nil {
  130. // Check if it's a timeout (expected)
  131. if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
  132. continue
  133. }
  134. // Check if server is stopping
  135. select {
  136. case <-s.stopChan:
  137. return
  138. default:
  139. }
  140. s.logger.Error("UDP read error", "err", err)
  141. continue
  142. }
  143. if n < 2 {
  144. s.logger.Debug("packet too short", "size", n, "addr", addr)
  145. continue
  146. }
  147. // Copy packet data to avoid buffer reuse issues
  148. packet := make([]byte, n)
  149. copy(packet, buf[:n])
  150. // Handle packet in goroutine to not block receive loop
  151. go s.handlePacket(addr, packet)
  152. }
  153. }
  154. // handlePacket processes a single incoming packet
  155. func (s *LegacyServer) handlePacket(addr *net.UDPAddr, packet []byte) {
  156. // Detect protocol version
  157. version, err := DetectProtocolVersion(packet)
  158. if err != nil {
  159. s.logger.Debug("failed to detect protocol version",
  160. "err", err,
  161. "addr", addr,
  162. )
  163. return
  164. }
  165. // Check if version is supported
  166. if !s.config.SupportsVersion(int(version)) {
  167. s.logger.Debug("unsupported protocol version",
  168. "version", version,
  169. "addr", addr,
  170. "size", len(packet),
  171. "hex", fmt.Sprintf("%X", packet),
  172. )
  173. return
  174. }
  175. // Get or create session
  176. session := s.sessions.GetSessionByAddr(addr)
  177. // Dispatch to appropriate handler
  178. if err := s.dispatcher.Dispatch(session, addr, packet); err != nil {
  179. s.logger.Debug("packet dispatch error",
  180. "err", err,
  181. "addr", addr,
  182. "version", version,
  183. )
  184. }
  185. }
  186. // SendPacket sends a packet to a specific address
  187. func (s *LegacyServer) SendPacket(addr *net.UDPAddr, packet []byte) error {
  188. s.mu.RLock()
  189. if !s.running || s.conn == nil {
  190. s.mu.RUnlock()
  191. return ErrServerClosed
  192. }
  193. conn := s.conn
  194. s.mu.RUnlock()
  195. s.logger.Debug("sending packet",
  196. "to", addr.String(),
  197. "size", len(packet),
  198. "hex", fmt.Sprintf("%X", packet),
  199. )
  200. _, err := conn.WriteToUDP(packet, addr)
  201. if err != nil {
  202. s.logger.Debug("failed to send packet",
  203. "err", err,
  204. "addr", addr,
  205. "size", len(packet),
  206. )
  207. return err
  208. }
  209. return nil
  210. }
  211. // SendToSession sends a packet to a session
  212. func (s *LegacyServer) SendToSession(session *LegacySession, packet []byte) error {
  213. if session == nil || session.Addr == nil {
  214. return errors.New("invalid session")
  215. }
  216. return s.SendPacket(session.Addr, packet)
  217. }
  218. // BroadcastToAll sends a packet to all connected sessions
  219. func (s *LegacyServer) BroadcastToAll(packet []byte) {
  220. sessions := s.sessions.GetAllSessions()
  221. for _, session := range sessions {
  222. if err := s.SendToSession(session, packet); err != nil {
  223. s.logger.Debug("broadcast send failed",
  224. "uin", session.UIN,
  225. "err", err,
  226. )
  227. }
  228. }
  229. }
  230. // IsRunning returns whether the server is currently running
  231. func (s *LegacyServer) IsRunning() bool {
  232. s.mu.RLock()
  233. defer s.mu.RUnlock()
  234. return s.running
  235. }
  236. // SessionCount returns the number of active sessions
  237. func (s *LegacyServer) SessionCount() int {
  238. return s.sessions.Count()
  239. }
  240. // ListenAndServe starts the server and blocks until it's stopped.
  241. // This method matches the interface used by other servers in the project.
  242. func (s *LegacyServer) ListenAndServe() error {
  243. ctx := context.Background()
  244. if err := s.Start(ctx); err != nil {
  245. return err
  246. }
  247. // Block until server is stopped
  248. <-s.stopChan
  249. return nil
  250. }
  251. // Shutdown gracefully shuts down the server.
  252. // This method matches the interface used by other servers in the project.
  253. func (s *LegacyServer) Shutdown(ctx context.Context) error {
  254. return s.Stop()
  255. }