webapi_session.go 12 KB

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