webapi_session.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454
  1. package state
  2. import (
  3. "context"
  4. "crypto/rand"
  5. "encoding/hex"
  6. "errors"
  7. "log/slog"
  8. "sync"
  9. "time"
  10. "github.com/mk6i/retro-aim-server/server/webapi/types"
  11. "github.com/mk6i/retro-aim-server/wire"
  12. )
  13. var (
  14. // ErrNoWebAPISession is returned when a WebAPI session is not found.
  15. ErrNoWebAPISession = errors.New("WebAPI session not found")
  16. // ErrWebAPISessionExpired is returned when a WebAPI session has expired.
  17. ErrWebAPISessionExpired = errors.New("WebAPI session expired")
  18. )
  19. // WebAPISession represents an active Web AIM API session.
  20. type WebAPISession struct {
  21. AimSID string // Unique session ID for web client
  22. ScreenName DisplayScreenName // User identity
  23. OSCARSession *Session // Bridge to existing OSCAR session
  24. Events []string // Subscribed event types
  25. EventQueue *types.EventQueue // Per-session event queue
  26. DevID string // Developer ID that created this session
  27. ClientName string // Client application name
  28. ClientVersion string // Client application version
  29. CreatedAt time.Time // Session creation time
  30. LastAccessed time.Time // Last activity time
  31. ExpiresAt time.Time // Session expiration time
  32. FetchTimeout int // Long-polling timeout in milliseconds
  33. TimeToNextFetch int // Suggested delay before next fetch
  34. RemoteAddr string // Client IP address
  35. logger *slog.Logger // Logger for debugging
  36. }
  37. // IsExpired checks if the session has expired.
  38. func (s *WebAPISession) IsExpired() bool {
  39. return time.Now().After(s.ExpiresAt)
  40. }
  41. // Touch updates the last accessed time and extends expiration if needed.
  42. func (s *WebAPISession) Touch() {
  43. s.LastAccessed = time.Now()
  44. // Extend expiration by 60 minutes from last access
  45. newExpiry := s.LastAccessed.Add(60 * time.Minute)
  46. if newExpiry.After(s.ExpiresAt) {
  47. s.ExpiresAt = newExpiry
  48. }
  49. }
  50. // IsSubscribedTo checks if the session is subscribed to a specific event type.
  51. func (s *WebAPISession) IsSubscribedTo(eventType string) bool {
  52. for _, event := range s.Events {
  53. if event == eventType {
  54. return true
  55. }
  56. }
  57. return false
  58. }
  59. // StartListeningToOSCARSession starts a goroutine that listens to the OSCAR session's
  60. // message channel and converts SNAC messages into WebAPI events.
  61. func (s *WebAPISession) StartListeningToOSCARSession() {
  62. if s.OSCARSession == nil {
  63. return
  64. }
  65. // Start goroutine to listen for OSCAR messages
  66. go func() {
  67. msgCh := s.OSCARSession.ReceiveMessage()
  68. for {
  69. select {
  70. case msg, ok := <-msgCh:
  71. if !ok {
  72. // Channel closed, OSCAR session ended
  73. return
  74. }
  75. s.handleSNACMessage(msg)
  76. case <-s.OSCARSession.Closed():
  77. // OSCAR session closed
  78. return
  79. }
  80. }
  81. }()
  82. }
  83. // handleSNACMessage converts a SNAC message into WebAPI events and pushes them to the event queue.
  84. func (s *WebAPISession) handleSNACMessage(msg wire.SNACMessage) {
  85. if s.EventQueue == nil {
  86. return
  87. }
  88. // Convert SNAC message to WebAPI events based on food group and subgroup
  89. switch msg.Frame.FoodGroup {
  90. case wire.ICBM:
  91. s.handleICBMMessage(msg)
  92. case wire.Buddy:
  93. s.handleBuddyMessage(msg)
  94. }
  95. }
  96. // handleICBMMessage handles ICBM (instant messaging) SNAC messages.
  97. func (s *WebAPISession) handleICBMMessage(msg wire.SNACMessage) {
  98. switch msg.Frame.SubGroup {
  99. case wire.ICBMChannelMsgToClient:
  100. s.handleIncomingIM(msg)
  101. case wire.ICBMClientEvent:
  102. s.handleTypingNotification(msg)
  103. }
  104. }
  105. // handleIncomingIM handles incoming instant messages.
  106. func (s *WebAPISession) handleIncomingIM(msg wire.SNACMessage) {
  107. if !s.IsSubscribedTo("im") {
  108. return
  109. }
  110. body, ok := msg.Body.(wire.SNAC_0x04_0x07_ICBMChannelMsgToClient)
  111. if !ok {
  112. return
  113. }
  114. // Extract message text from TLV data
  115. var messageText string
  116. if msgData, hasMsg := body.TLVRestBlock.Bytes(wire.ICBMTLVAOLIMData); hasMsg {
  117. if text, err := wire.UnmarshalICBMMessageText(msgData); err == nil {
  118. messageText = text
  119. }
  120. }
  121. if messageText == "" {
  122. return
  123. }
  124. // Check if it's an auto-response (channel 2)
  125. autoResponse := body.ChannelID == 0x0002
  126. // Create IM event
  127. imEvent := types.IMEvent{
  128. From: body.ScreenName,
  129. Message: messageText,
  130. Timestamp: float64(time.Now().Unix()),
  131. AutoResp: autoResponse,
  132. }
  133. s.EventQueue.Push(types.EventTypeIM, imEvent)
  134. }
  135. // handleTypingNotification handles typing notifications.
  136. func (s *WebAPISession) handleTypingNotification(msg wire.SNACMessage) {
  137. if !s.IsSubscribedTo("typing") {
  138. return
  139. }
  140. body, ok := msg.Body.(wire.SNAC_0x04_0x14_ICBMClientEvent)
  141. if !ok {
  142. return
  143. }
  144. // Event types: 0=stopped typing, 1=text typed, 2=typing
  145. isTyping := body.Event == 1 || body.Event == 2
  146. typingEvent := types.TypingEvent{
  147. From: body.ScreenName,
  148. Typing: isTyping,
  149. }
  150. s.EventQueue.Push(types.EventTypeTyping, typingEvent)
  151. }
  152. // handleBuddyMessage handles buddy/presence SNAC messages.
  153. func (s *WebAPISession) handleBuddyMessage(msg wire.SNACMessage) {
  154. switch msg.Frame.SubGroup {
  155. case wire.BuddyArrived:
  156. s.handleBuddyArrived(msg)
  157. case wire.BuddyDeparted:
  158. s.handleBuddyDeparted(msg)
  159. }
  160. }
  161. // handleBuddyArrived handles when a buddy comes online.
  162. func (s *WebAPISession) handleBuddyArrived(msg wire.SNACMessage) {
  163. if !s.IsSubscribedTo("presence") {
  164. return
  165. }
  166. body, ok := msg.Body.(wire.SNAC_0x03_0x0B_BuddyArrived)
  167. if !ok {
  168. return
  169. }
  170. presenceEvent := types.PresenceEvent{
  171. AimID: body.ScreenName,
  172. State: "online",
  173. UserType: "aim",
  174. }
  175. s.EventQueue.Push(types.EventTypePresence, presenceEvent)
  176. }
  177. // handleBuddyDeparted handles when a buddy goes offline.
  178. func (s *WebAPISession) handleBuddyDeparted(msg wire.SNACMessage) {
  179. if !s.IsSubscribedTo("presence") {
  180. return
  181. }
  182. body, ok := msg.Body.(wire.SNAC_0x03_0x0C_BuddyDeparted)
  183. if !ok {
  184. return
  185. }
  186. presenceEvent := types.PresenceEvent{
  187. AimID: body.ScreenName,
  188. State: "offline",
  189. UserType: "aim",
  190. }
  191. s.EventQueue.Push(types.EventTypePresence, presenceEvent)
  192. }
  193. // WebAPISessionManager manages Web API sessions with thread-safe operations.
  194. type WebAPISessionManager struct {
  195. sessions map[string]*WebAPISession // Keyed by aimsid
  196. byUser map[IdentScreenName]*WebAPISession // Keyed by screen name
  197. mu sync.RWMutex
  198. cleanupTicker *time.Ticker
  199. stopCleanup chan struct{}
  200. }
  201. // NewWebAPISessionManager creates a new WebAPI session manager.
  202. func NewWebAPISessionManager() *WebAPISessionManager {
  203. mgr := &WebAPISessionManager{
  204. sessions: make(map[string]*WebAPISession),
  205. byUser: make(map[IdentScreenName]*WebAPISession),
  206. stopCleanup: make(chan struct{}),
  207. }
  208. // Start cleanup goroutine to remove expired sessions
  209. mgr.cleanupTicker = time.NewTicker(1 * time.Minute)
  210. go mgr.cleanupExpiredSessions()
  211. return mgr
  212. }
  213. // CreateSession creates a new WebAPI session.
  214. func (m *WebAPISessionManager) CreateSession(ctx context.Context, screenName DisplayScreenName, devID string, events []string, oscarSession *Session, logger *slog.Logger) (*WebAPISession, error) {
  215. m.mu.Lock()
  216. defer m.mu.Unlock()
  217. // Check if user already has an active session
  218. identName := screenName.IdentScreenName()
  219. if existing, exists := m.byUser[identName]; exists {
  220. // Remove the old session
  221. delete(m.sessions, existing.AimSID)
  222. }
  223. // Generate unique session ID
  224. aimsid, err := generateSessionID()
  225. if err != nil {
  226. return nil, err
  227. }
  228. now := time.Now()
  229. session := &WebAPISession{
  230. AimSID: aimsid,
  231. ScreenName: screenName,
  232. OSCARSession: oscarSession,
  233. Events: events,
  234. EventQueue: types.NewEventQueue(1000), // Max 1000 events per session
  235. DevID: devID,
  236. CreatedAt: now,
  237. LastAccessed: now,
  238. ExpiresAt: now.Add(60 * time.Minute), // 60 minute initial expiry
  239. FetchTimeout: 60000, // 60 seconds default for better stability
  240. TimeToNextFetch: 500, // 500ms suggested delay
  241. logger: logger,
  242. }
  243. m.sessions[aimsid] = session
  244. m.byUser[identName] = session
  245. // Start listening to OSCAR session message channel
  246. session.StartListeningToOSCARSession()
  247. return session, nil
  248. }
  249. // GetSession retrieves a session by aimsid.
  250. func (m *WebAPISessionManager) GetSession(ctx context.Context, aimsid string) (*WebAPISession, error) {
  251. m.mu.RLock()
  252. defer m.mu.RUnlock()
  253. session, exists := m.sessions[aimsid]
  254. if !exists {
  255. return nil, ErrNoWebAPISession
  256. }
  257. if session.IsExpired() {
  258. return nil, ErrWebAPISessionExpired
  259. }
  260. return session, nil
  261. }
  262. // GetSessionByUser retrieves a session by screen name.
  263. func (m *WebAPISessionManager) GetSessionByUser(ctx context.Context, screenName IdentScreenName) (*WebAPISession, error) {
  264. m.mu.RLock()
  265. defer m.mu.RUnlock()
  266. session, exists := m.byUser[screenName]
  267. if !exists {
  268. return nil, ErrNoWebAPISession
  269. }
  270. if session.IsExpired() {
  271. return nil, ErrWebAPISessionExpired
  272. }
  273. return session, nil
  274. }
  275. // RemoveSession removes a session by aimsid.
  276. func (m *WebAPISessionManager) RemoveSession(ctx context.Context, aimsid string) error {
  277. m.mu.Lock()
  278. defer m.mu.Unlock()
  279. session, exists := m.sessions[aimsid]
  280. if !exists {
  281. return ErrNoWebAPISession
  282. }
  283. delete(m.sessions, aimsid)
  284. delete(m.byUser, session.ScreenName.IdentScreenName())
  285. // Close the event queue to unblock any waiting fetches
  286. if session.EventQueue != nil {
  287. session.EventQueue.Close()
  288. }
  289. return nil
  290. }
  291. // TouchSession updates the last accessed time for a session.
  292. func (m *WebAPISessionManager) TouchSession(ctx context.Context, aimsid string) error {
  293. m.mu.Lock()
  294. defer m.mu.Unlock()
  295. session, exists := m.sessions[aimsid]
  296. if !exists {
  297. return ErrNoWebAPISession
  298. }
  299. session.Touch()
  300. return nil
  301. }
  302. // GetAllSessions returns all active sessions (for monitoring/admin).
  303. func (m *WebAPISessionManager) GetAllSessions(ctx context.Context) []*WebAPISession {
  304. m.mu.RLock()
  305. defer m.mu.RUnlock()
  306. sessions := make([]*WebAPISession, 0, len(m.sessions))
  307. for _, session := range m.sessions {
  308. if !session.IsExpired() {
  309. sessions = append(sessions, session)
  310. }
  311. }
  312. return sessions
  313. }
  314. // GetSessionsByScreenName returns all sessions for a given screen name.
  315. func (m *WebAPISessionManager) GetSessionsByScreenName(ctx context.Context, screenName DisplayScreenName) []*WebAPISession {
  316. m.mu.RLock()
  317. defer m.mu.RUnlock()
  318. var sessions []*WebAPISession
  319. identScreenName := screenName.IdentScreenName()
  320. // Check both the byUser map and iterate through all sessions
  321. // since a user might have multiple sessions
  322. for _, session := range m.sessions {
  323. if session.ScreenName.IdentScreenName() == identScreenName {
  324. sessions = append(sessions, session)
  325. }
  326. }
  327. return sessions
  328. }
  329. // cleanupExpiredSessions periodically removes expired sessions.
  330. func (m *WebAPISessionManager) cleanupExpiredSessions() {
  331. for {
  332. select {
  333. case <-m.cleanupTicker.C:
  334. m.mu.Lock()
  335. now := time.Now()
  336. var toRemove []string
  337. for aimsid, session := range m.sessions {
  338. if now.After(session.ExpiresAt) {
  339. toRemove = append(toRemove, aimsid)
  340. }
  341. }
  342. for _, aimsid := range toRemove {
  343. session := m.sessions[aimsid]
  344. delete(m.sessions, aimsid)
  345. delete(m.byUser, session.ScreenName.IdentScreenName())
  346. if session.EventQueue != nil {
  347. session.EventQueue.Close()
  348. }
  349. }
  350. m.mu.Unlock()
  351. case <-m.stopCleanup:
  352. m.cleanupTicker.Stop()
  353. return
  354. }
  355. }
  356. }
  357. // Shutdown stops the session manager and cleans up resources.
  358. func (m *WebAPISessionManager) Shutdown(ctx context.Context) {
  359. close(m.stopCleanup)
  360. m.mu.Lock()
  361. defer m.mu.Unlock()
  362. // Close all event queues
  363. for _, session := range m.sessions {
  364. if session.EventQueue != nil {
  365. session.EventQueue.Close()
  366. }
  367. }
  368. // Clear all sessions
  369. m.sessions = make(map[string]*WebAPISession)
  370. m.byUser = make(map[IdentScreenName]*WebAPISession)
  371. }
  372. // generateSessionID creates a cryptographically secure session ID.
  373. func generateSessionID() (string, error) {
  374. bytes := make([]byte, 32) // 256 bits
  375. if _, err := rand.Read(bytes); err != nil {
  376. return "", err
  377. }
  378. return hex.EncodeToString(bytes), nil
  379. }