connection.go 7.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285
  1. package server
  2. import (
  3. "bytes"
  4. "context"
  5. "errors"
  6. "github.com/mkaminski/goaim/user"
  7. "io"
  8. "log"
  9. "log/slog"
  10. "net"
  11. "os"
  12. "github.com/google/uuid"
  13. "github.com/mkaminski/goaim/oscar"
  14. )
  15. var (
  16. ErrUnsupportedSubGroup = errors.New("unimplemented subgroup, your client version may be unsupported")
  17. )
  18. type (
  19. incomingMessage struct {
  20. flap oscar.FlapFrame
  21. payload *bytes.Buffer
  22. }
  23. alertHandler func(ctx context.Context, msg oscar.XMessage, w io.Writer, u *uint32) error
  24. clientReqHandler func(ctx context.Context, r io.Reader, w io.Writer, u *uint32) error
  25. )
  26. func consumeFLAPFrames(r io.Reader, msgCh chan incomingMessage, errCh chan error) {
  27. defer close(msgCh)
  28. defer close(errCh)
  29. for {
  30. in := incomingMessage{}
  31. if err := oscar.Unmarshal(&in.flap, r); err != nil {
  32. errCh <- err
  33. return
  34. }
  35. if in.flap.FrameType == oscar.FlapFrameData {
  36. buf := make([]byte, in.flap.PayloadLength)
  37. if _, err := r.Read(buf); err != nil {
  38. errCh <- err
  39. return
  40. }
  41. in.payload = bytes.NewBuffer(buf)
  42. }
  43. msgCh <- in
  44. }
  45. }
  46. func dispatchIncomingMessages(ctx context.Context, sess *user.Session, seq uint32, rw io.ReadWriter, logger *slog.Logger, fn clientReqHandler, alertHandler alertHandler) {
  47. // buffered so that the go routine has room to exit
  48. msgCh := make(chan incomingMessage, 1)
  49. readErrCh := make(chan error, 1)
  50. go consumeFLAPFrames(rw, msgCh, readErrCh)
  51. for {
  52. select {
  53. case m := <-msgCh:
  54. switch m.flap.FrameType {
  55. case oscar.FlapFrameData:
  56. // route a client request to the appropriate service handler. the
  57. // handler may write a response to the client connection.
  58. if err := fn(ctx, m.payload, rw, &seq); err != nil {
  59. return
  60. }
  61. case oscar.FlapFrameSignon:
  62. logger.ErrorContext(ctx, "shouldn't get FlapFrameSignon", "flap", m.flap)
  63. case oscar.FlapFrameError:
  64. logger.ErrorContext(ctx, "got FlapFrameError", "flap", m.flap)
  65. return
  66. case oscar.FlapFrameSignoff:
  67. logger.InfoContext(ctx, "got FlapFrameSignoff", "flap", m.flap)
  68. return
  69. case oscar.FlapFrameKeepAlive:
  70. logger.DebugContext(ctx, "keepalive heartbeat")
  71. default:
  72. logger.ErrorContext(ctx, "got unknown FLAP frame type", "flap", m.flap)
  73. return
  74. }
  75. case m := <-sess.RecvMessage():
  76. // forward a notification sent from another client to this client
  77. if err := alertHandler(ctx, m, rw, &seq); err != nil {
  78. logRequestError(ctx, logger, m.SnacFrame, err)
  79. return
  80. }
  81. logRequest(ctx, logger, m.SnacFrame, m.SnacOut)
  82. case <-sess.Closed():
  83. // gracefully disconnect so that the client does not try to
  84. // reconnect when the connection closes.
  85. flap := oscar.FlapFrame{
  86. StartMarker: 42,
  87. FrameType: oscar.FlapFrameSignoff,
  88. Sequence: uint16(seq),
  89. PayloadLength: uint16(0),
  90. }
  91. if err := oscar.Marshal(flap, rw); err != nil {
  92. logger.ErrorContext(ctx, "unable to gracefully disconnect user", "err", err)
  93. }
  94. return
  95. case err := <-readErrCh:
  96. // handle a read error
  97. switch {
  98. case errors.Is(io.EOF, err):
  99. fallthrough
  100. case errors.Is(ErrSignedOff, err):
  101. logger.InfoContext(ctx, "client signed off")
  102. default:
  103. logger.ErrorContext(ctx, "client disconnected with error", "err", err)
  104. }
  105. return
  106. }
  107. }
  108. }
  109. func HandleChatConnection(ctx context.Context, cr *ChatRegistry, rw io.ReadWriter, router ChatServiceRouter, logger *slog.Logger) {
  110. cookie, seq, err := VerifyChatLogin(rw)
  111. if err != nil {
  112. logger.ErrorContext(ctx, "user disconnected with error", "err", err.Error())
  113. return
  114. }
  115. room, err := cr.Retrieve(string(cookie.Cookie))
  116. if err != nil {
  117. logger.ErrorContext(ctx, "unable to find chat room", "err", err.Error())
  118. return
  119. }
  120. chatSess, found := room.Retrieve(cookie.SessID)
  121. if !found {
  122. logger.ErrorContext(ctx, "unable to find user for session", "sessID", cookie.SessID)
  123. return
  124. }
  125. defer chatSess.Close()
  126. go func() {
  127. <-chatSess.Closed()
  128. AlertUserLeft(ctx, chatSess, room)
  129. room.Remove(chatSess)
  130. cr.MaybeRemoveRoom(room.Cookie)
  131. }()
  132. ctx = context.WithValue(ctx, "screenName", chatSess.ScreenName())
  133. if err := router.WriteOServiceHostOnline(rw, &seq); err != nil {
  134. logger.ErrorContext(ctx, "error WriteOServiceHostOnline")
  135. }
  136. fnClientReqHandler := func(ctx context.Context, r io.Reader, w io.Writer, seq *uint32) error {
  137. return router.Route(ctx, chatSess, r, w, seq, room)
  138. }
  139. fnAlertHandler := func(ctx context.Context, msg oscar.XMessage, w io.Writer, seq *uint32) error {
  140. return writeOutSNAC(oscar.SnacFrame{}, msg.SnacFrame, msg.SnacOut, seq, w)
  141. }
  142. dispatchIncomingMessages(ctx, chatSess, seq, rw, logger, fnClientReqHandler, fnAlertHandler)
  143. }
  144. func HandleAuthConnection(cfg Config, sm *InMemorySessionManager, fm *FeedbagStore, conn net.Conn) {
  145. defer conn.Close()
  146. seq := uint32(100)
  147. _, err := SendAndReceiveSignonFrame(conn, &seq)
  148. if err != nil {
  149. log.Println(err)
  150. return
  151. }
  152. err = ReceiveAndSendAuthChallenge(cfg, fm, conn, conn, &seq, uuid.New)
  153. if err != nil {
  154. log.Println(err)
  155. return
  156. }
  157. err = ReceiveAndSendBUCPLoginRequest(cfg, sm, fm, conn, conn, &seq, uuid.New)
  158. if err != nil {
  159. log.Println(err)
  160. return
  161. }
  162. }
  163. func HandleBOSConnection(ctx context.Context, conn net.Conn, router BOSServiceRouter, logger *slog.Logger) {
  164. sess, seq, err := router.VerifyLogin(conn)
  165. if err != nil {
  166. logger.ErrorContext(ctx, "user disconnected with error", "err", err.Error())
  167. return
  168. }
  169. defer sess.Close()
  170. defer conn.Close()
  171. go func() {
  172. <-sess.Closed()
  173. router.Signout(ctx, logger, sess)
  174. }()
  175. ctx = context.WithValue(ctx, "screenName", sess.ScreenName())
  176. if err := router.WriteOServiceHostOnline(conn, &seq); err != nil {
  177. logger.ErrorContext(ctx, "error WriteOServiceHostOnline")
  178. }
  179. fnClientReqHandler := func(ctx context.Context, r io.Reader, w io.Writer, seq *uint32) error {
  180. return router.Route(ctx, sess, r, w, seq)
  181. }
  182. fnAlertHandler := func(ctx context.Context, msg oscar.XMessage, w io.Writer, seq *uint32) error {
  183. return writeOutSNAC(oscar.SnacFrame{}, msg.SnacFrame, msg.SnacOut, seq, w)
  184. }
  185. dispatchIncomingMessages(ctx, sess, seq, conn, logger, fnClientReqHandler, fnAlertHandler)
  186. }
  187. func ListenChat(cfg Config, router ChatServiceRouter, cr *ChatRegistry, logger *slog.Logger) {
  188. addr := Address("", cfg.ChatPort)
  189. listener, err := net.Listen("tcp", addr)
  190. if err != nil {
  191. logger.Error("unable to bind chat server address", "err", err.Error())
  192. os.Exit(1)
  193. }
  194. defer listener.Close()
  195. logger.Info("starting service", "addr", addr)
  196. for {
  197. conn, err := listener.Accept()
  198. if err != nil {
  199. log.Println(err)
  200. continue
  201. }
  202. ctx := context.Background()
  203. ctx = context.WithValue(ctx, "ip", conn.RemoteAddr().String())
  204. logger.DebugContext(ctx, "accepted connection")
  205. go func() {
  206. HandleChatConnection(ctx, cr, conn, router, logger)
  207. conn.Close()
  208. }()
  209. }
  210. }
  211. func ListenBOS(cfg Config, router BOSServiceRouter, logger *slog.Logger) {
  212. addr := Address("", cfg.BOSPort)
  213. listener, err := net.Listen("tcp", addr)
  214. if err != nil {
  215. logger.Error("unable to bind BOS server address", "err", err.Error())
  216. os.Exit(1)
  217. }
  218. defer listener.Close()
  219. logger.Info("starting service", "addr", addr)
  220. for {
  221. conn, err := listener.Accept()
  222. if err != nil {
  223. log.Println(err)
  224. continue
  225. }
  226. ctx := context.Background()
  227. ctx = context.WithValue(ctx, "ip", conn.RemoteAddr().String())
  228. logger.DebugContext(ctx, "accepted connection")
  229. go HandleBOSConnection(ctx, conn, router, logger)
  230. }
  231. }
  232. func ListenBUCPLogin(cfg Config, err error, logger *slog.Logger, sm *InMemorySessionManager, fm *FeedbagStore) {
  233. addr := Address("", cfg.OSCARPort)
  234. listener, err := net.Listen("tcp", addr)
  235. if err != nil {
  236. logger.Error("unable to bind OSCAR server address", "err", err.Error())
  237. os.Exit(1)
  238. }
  239. defer listener.Close()
  240. logger.Info("starting OSCAR server", "addr", addr)
  241. for {
  242. conn, err := listener.Accept()
  243. if err != nil {
  244. log.Println(err)
  245. continue
  246. }
  247. go HandleAuthConnection(cfg, sm, fm, conn)
  248. }
  249. }