webapi_session.go 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548
  1. package state
  2. import (
  3. "context"
  4. "crypto/rand"
  5. "encoding/hex"
  6. "errors"
  7. "log/slog"
  8. mrand "math/rand/v2"
  9. "strconv"
  10. "sync"
  11. "time"
  12. "github.com/mk6i/open-oscar-server/server/webapi/types"
  13. "github.com/mk6i/open-oscar-server/wire"
  14. )
  15. var (
  16. // ErrNoWebAPISession is returned when a WebAPI session is not found.
  17. ErrNoWebAPISession = errors.New("WebAPI session not found")
  18. // ErrWebAPISessionExpired is returned when a WebAPI session has expired.
  19. ErrWebAPISessionExpired = errors.New("WebAPI session expired")
  20. )
  21. // WebAPISession represents an active Web AIM API session.
  22. type WebAPISession struct {
  23. AimSID string // Unique session ID for web client
  24. ScreenName DisplayScreenName // User identity
  25. OSCARSession *SessionInstance // Bridge to existing OSCAR session
  26. Events []string // Subscribed event types
  27. EventQueue *types.EventQueue // Per-session event queue
  28. DevID string // Developer ID that created this session
  29. ClientName string // Client application name
  30. ClientVersion string // Client application version
  31. CreatedAt time.Time // SessionInstance creation time
  32. LastAccessed time.Time // Last activity time
  33. ExpiresAt time.Time // SessionInstance expiration time
  34. FetchTimeout int // Long-polling timeout in milliseconds
  35. TimeToNextFetch int // Suggested delay before next fetch
  36. RemoteAddr string // Client IP address
  37. TempBuddies map[string]bool // Temporary buddies for this session only
  38. BuddyListRefresher func(ctx context.Context) (interface{}, error) // Called on feedbag changes to push buddylist event
  39. PermitDenyRefresher func(ctx context.Context) (interface{}, error) // Called on feedbag changes to push permitDeny event
  40. imLog map[string][]WebAPIStoredIM
  41. imLogMu sync.Mutex
  42. logger *slog.Logger // Logger for debugging
  43. }
  44. // IsExpired checks if the session has expired.
  45. func (s *WebAPISession) IsExpired() bool {
  46. return time.Now().After(s.ExpiresAt)
  47. }
  48. // Touch updates the last accessed time and extends expiration if needed.
  49. func (s *WebAPISession) Touch() {
  50. s.LastAccessed = time.Now()
  51. // Extend expiration by 60 minutes from last access
  52. newExpiry := s.LastAccessed.Add(60 * time.Minute)
  53. if newExpiry.After(s.ExpiresAt) {
  54. s.ExpiresAt = newExpiry
  55. }
  56. }
  57. // IsSubscribedTo checks if the session is subscribed to a specific event type.
  58. func (s *WebAPISession) IsSubscribedTo(eventType string) bool {
  59. for _, event := range s.Events {
  60. if event == eventType {
  61. return true
  62. }
  63. }
  64. return false
  65. }
  66. // StartListeningToOSCARSession starts a goroutine that listens to the OSCAR session's
  67. // message channel and converts SNAC messages into WebAPI events.
  68. func (s *WebAPISession) StartListeningToOSCARSession() {
  69. if s.OSCARSession == nil {
  70. return
  71. }
  72. // Start goroutine to listen for OSCAR messages
  73. go func() {
  74. msgCh := s.OSCARSession.ReceiveMessage()
  75. for {
  76. select {
  77. case msg, ok := <-msgCh:
  78. if !ok {
  79. // Channel closed, OSCAR session ended
  80. return
  81. }
  82. s.handleSNACMessage(msg)
  83. case <-s.OSCARSession.Closed():
  84. // OSCAR session closed
  85. return
  86. }
  87. }
  88. }()
  89. }
  90. // handleSNACMessage converts a SNAC message into WebAPI events and pushes them to the event queue.
  91. func (s *WebAPISession) handleSNACMessage(msg wire.SNACMessage) {
  92. if s.EventQueue == nil {
  93. return
  94. }
  95. // Convert SNAC message to WebAPI events based on food group and subgroup
  96. switch msg.Frame.FoodGroup {
  97. case wire.ICBM:
  98. s.handleICBMMessage(msg)
  99. case wire.Buddy:
  100. s.handleBuddyMessage(msg)
  101. case wire.Feedbag:
  102. s.handleFeedbagMessage(msg)
  103. }
  104. }
  105. // handleICBMMessage handles ICBM (instant messaging) SNAC messages.
  106. func (s *WebAPISession) handleICBMMessage(msg wire.SNACMessage) {
  107. switch msg.Frame.SubGroup {
  108. case wire.ICBMChannelMsgToClient:
  109. s.handleIncomingIM(msg)
  110. case wire.ICBMClientEvent:
  111. s.handleTypingNotification(msg)
  112. }
  113. }
  114. // handleIncomingIM handles incoming instant messages.
  115. func (s *WebAPISession) handleIncomingIM(msg wire.SNACMessage) {
  116. if !s.IsSubscribedTo("im") {
  117. return
  118. }
  119. body, ok := msg.Body.(wire.SNAC_0x04_0x07_ICBMChannelMsgToClient)
  120. if !ok {
  121. return
  122. }
  123. // Extract message text from TLV data
  124. var messageText string
  125. if msgData, hasMsg := body.Bytes(wire.ICBMTLVAOLIMData); hasMsg {
  126. if text, err := wire.UnmarshalICBMMessageText(msgData); err == nil {
  127. messageText = text
  128. }
  129. }
  130. if messageText == "" {
  131. return
  132. }
  133. // Check if it's an auto-response (channel 2)
  134. autoResponse := body.ChannelID == 0x0002
  135. // msgId must be unique per delivered event. The OSCAR cookie is not a
  136. // reliable unique id (some clients reuse it across messages), and the web
  137. // client dedupes its conversation list by msgId, silently dropping any
  138. // collisions. Mint a fresh random id instead of reusing body.Cookie.
  139. msgID := strconv.FormatUint(mrand.Uint64(), 16)
  140. partner := body.ScreenName
  141. nowSec := time.Now().Unix()
  142. s.AddStoredIM(partner, partner, messageText, msgID, nowSec)
  143. // Create IM event
  144. imEvent := types.IMEvent{
  145. Source: types.UserInfo{
  146. AimID: body.ScreenName,
  147. DisplayID: body.ScreenName,
  148. UserType: "aim",
  149. State: "online",
  150. },
  151. Message: messageText,
  152. MsgID: msgID,
  153. Timestamp: float64(time.Now().Unix()),
  154. AutoResp: autoResponse,
  155. }
  156. s.EventQueue.Push(types.EventTypeIM, imEvent)
  157. if s.IsSubscribedTo("conversation") {
  158. s.EventQueue.Push(types.EventTypeConversation, types.ConversationEventData("update", []map[string]interface{}{
  159. types.ConversationEntry(
  160. body.ScreenName,
  161. body.ScreenName,
  162. messageText,
  163. msgID,
  164. body.ScreenName,
  165. false,
  166. 1,
  167. ),
  168. }))
  169. }
  170. }
  171. // handleTypingNotification handles typing notifications.
  172. func (s *WebAPISession) handleTypingNotification(msg wire.SNACMessage) {
  173. if !s.IsSubscribedTo("typing") {
  174. return
  175. }
  176. body, ok := msg.Body.(wire.SNAC_0x04_0x14_ICBMClientEvent)
  177. if !ok {
  178. return
  179. }
  180. // Event types: 0x0000=none, 0x0001=typed (paused), 0x0002=typing
  181. var typingStatus string
  182. switch body.Event {
  183. case 0x0002:
  184. typingStatus = "typing"
  185. case 0x0001:
  186. typingStatus = "typed"
  187. default:
  188. typingStatus = "none"
  189. }
  190. typingEvent := types.TypingEvent{
  191. AimID: body.ScreenName,
  192. TypingStatus: typingStatus,
  193. }
  194. s.EventQueue.Push(types.EventTypeTyping, typingEvent)
  195. }
  196. // handleBuddyMessage handles buddy/presence SNAC messages.
  197. func (s *WebAPISession) handleBuddyMessage(msg wire.SNACMessage) {
  198. switch msg.Frame.SubGroup {
  199. case wire.BuddyArrived:
  200. s.handleBuddyArrived(msg)
  201. case wire.BuddyDeparted:
  202. s.handleBuddyDeparted(msg)
  203. }
  204. }
  205. // handleBuddyArrived handles when a buddy comes online.
  206. func (s *WebAPISession) handleBuddyArrived(msg wire.SNACMessage) {
  207. if !s.IsSubscribedTo("presence") {
  208. return
  209. }
  210. body, ok := msg.Body.(wire.SNAC_0x03_0x0B_BuddyArrived)
  211. if !ok {
  212. return
  213. }
  214. stateStr := "online"
  215. // For BuddyArrived updates, infer presence state from the TLVUserInfo.
  216. // Away and invisible transitions are typically broadcast using BuddyArrived
  217. // with updated user flags/status bits, not BuddyDeparted.
  218. if body.IsInvisible() {
  219. stateStr = "offline"
  220. } else if body.IsAway() {
  221. stateStr = "away"
  222. } else if mask, ok := body.Uint32BE(wire.OServiceUserInfoStatus); ok {
  223. if mask&wire.OServiceUserStatusDND == wire.OServiceUserStatusDND {
  224. stateStr = "dnd"
  225. } else if mask&wire.OServiceUserStatusAway == wire.OServiceUserStatusAway {
  226. stateStr = "away"
  227. }
  228. }
  229. presenceEvent := types.PresenceEvent{
  230. AimID: body.ScreenName,
  231. State: stateStr,
  232. UserType: "aim",
  233. }
  234. s.EventQueue.Push(types.EventTypePresence, presenceEvent)
  235. }
  236. // handleBuddyDeparted handles when a buddy goes offline.
  237. func (s *WebAPISession) handleBuddyDeparted(msg wire.SNACMessage) {
  238. if !s.IsSubscribedTo("presence") {
  239. return
  240. }
  241. body, ok := msg.Body.(wire.SNAC_0x03_0x0C_BuddyDeparted)
  242. if !ok {
  243. return
  244. }
  245. presenceEvent := types.PresenceEvent{
  246. AimID: body.ScreenName,
  247. State: "offline",
  248. UserType: "aim",
  249. }
  250. s.EventQueue.Push(types.EventTypePresence, presenceEvent)
  251. }
  252. func (s *WebAPISession) handleFeedbagMessage(msg wire.SNACMessage) {
  253. switch msg.Frame.SubGroup {
  254. case wire.FeedbagInsertItem, wire.FeedbagUpdateItem, wire.FeedbagDeleteItem:
  255. if s.BuddyListRefresher != nil {
  256. groups, err := s.BuddyListRefresher(context.Background())
  257. if err != nil {
  258. s.logger.Error("failed to refresh buddy list after feedbag change", "err", err)
  259. } else {
  260. s.EventQueue.Push(types.EventTypeBuddyList, map[string]interface{}{"groups": groups})
  261. }
  262. }
  263. if msg.Frame.SubGroup == wire.FeedbagUpdateItem && s.PermitDenyRefresher != nil {
  264. body, ok := msg.Body.(wire.SNAC_0x13_0x09_FeedbagUpdateItem)
  265. if ok {
  266. for _, item := range body.Items {
  267. if item.ClassID == wire.FeedbagClassIDPermit ||
  268. item.ClassID == wire.FeedbagClassIDDeny ||
  269. item.ClassID == wire.FeedbagClassIdPdinfo {
  270. pdd, err := s.PermitDenyRefresher(context.Background())
  271. if err != nil {
  272. s.logger.Error("failed to refresh permit/deny after feedbag change", "err", err)
  273. } else {
  274. s.EventQueue.Push(types.EventTypePermitDeny, pdd)
  275. }
  276. break
  277. }
  278. }
  279. }
  280. }
  281. }
  282. }
  283. // WebAPISessionManager manages Web API sessions with thread-safe operations.
  284. type WebAPISessionManager struct {
  285. sessions map[string]*WebAPISession // Keyed by aimsid
  286. byUser map[IdentScreenName]*WebAPISession // Keyed by screen name
  287. mu sync.RWMutex
  288. cleanupTicker *time.Ticker
  289. stopCleanup chan struct{}
  290. }
  291. // NewWebAPISessionManager creates a new WebAPI session manager.
  292. func NewWebAPISessionManager() *WebAPISessionManager {
  293. mgr := &WebAPISessionManager{
  294. sessions: make(map[string]*WebAPISession),
  295. byUser: make(map[IdentScreenName]*WebAPISession),
  296. stopCleanup: make(chan struct{}),
  297. }
  298. // Start cleanup goroutine to remove expired sessions
  299. mgr.cleanupTicker = time.NewTicker(1 * time.Minute)
  300. go mgr.cleanupExpiredSessions()
  301. return mgr
  302. }
  303. // CreateSession creates a new WebAPI session.
  304. func (m *WebAPISessionManager) CreateSession(ctx context.Context, screenName DisplayScreenName, devID string, events []string, oscarSession *SessionInstance, logger *slog.Logger) (*WebAPISession, error) {
  305. m.mu.Lock()
  306. defer m.mu.Unlock()
  307. // Check if user already has an active session
  308. identName := screenName.IdentScreenName()
  309. if existing, exists := m.byUser[identName]; exists {
  310. // Remove the old session
  311. delete(m.sessions, existing.AimSID)
  312. }
  313. // Generate unique session ID
  314. aimsid, err := generateSessionID()
  315. if err != nil {
  316. return nil, err
  317. }
  318. now := time.Now()
  319. session := &WebAPISession{
  320. AimSID: aimsid,
  321. ScreenName: screenName,
  322. OSCARSession: oscarSession,
  323. Events: events,
  324. EventQueue: types.NewEventQueue(1000), // Max 1000 events per session
  325. DevID: devID,
  326. CreatedAt: now,
  327. LastAccessed: now,
  328. ExpiresAt: now.Add(60 * time.Minute), // 60 minute initial expiry
  329. FetchTimeout: 60000, // 60 seconds default for better stability
  330. TimeToNextFetch: 500, // 500ms suggested delay
  331. logger: logger,
  332. }
  333. m.sessions[aimsid] = session
  334. m.byUser[identName] = session
  335. // Start listening to OSCAR session message channel
  336. session.StartListeningToOSCARSession()
  337. return session, nil
  338. }
  339. // GetSession retrieves a session by aimsid.
  340. func (m *WebAPISessionManager) GetSession(ctx context.Context, aimsid string) (*WebAPISession, error) {
  341. m.mu.RLock()
  342. defer m.mu.RUnlock()
  343. session, exists := m.sessions[aimsid]
  344. if !exists {
  345. return nil, ErrNoWebAPISession
  346. }
  347. if session.IsExpired() {
  348. return nil, ErrWebAPISessionExpired
  349. }
  350. return session, nil
  351. }
  352. // GetSessionByUser retrieves a session by screen name.
  353. func (m *WebAPISessionManager) GetSessionByUser(ctx context.Context, screenName IdentScreenName) (*WebAPISession, error) {
  354. m.mu.RLock()
  355. defer m.mu.RUnlock()
  356. session, exists := m.byUser[screenName]
  357. if !exists {
  358. return nil, ErrNoWebAPISession
  359. }
  360. if session.IsExpired() {
  361. return nil, ErrWebAPISessionExpired
  362. }
  363. return session, nil
  364. }
  365. // RemoveSession removes a session by aimsid.
  366. func (m *WebAPISessionManager) RemoveSession(ctx context.Context, aimsid string) error {
  367. m.mu.Lock()
  368. defer m.mu.Unlock()
  369. session, exists := m.sessions[aimsid]
  370. if !exists {
  371. return ErrNoWebAPISession
  372. }
  373. delete(m.sessions, aimsid)
  374. delete(m.byUser, session.ScreenName.IdentScreenName())
  375. // CloseSession the event queue to unblock any waiting fetches
  376. if session.EventQueue != nil {
  377. session.EventQueue.Close()
  378. }
  379. return nil
  380. }
  381. // TouchSession updates the last accessed time for a session.
  382. func (m *WebAPISessionManager) TouchSession(ctx context.Context, aimsid string) error {
  383. m.mu.Lock()
  384. defer m.mu.Unlock()
  385. session, exists := m.sessions[aimsid]
  386. if !exists {
  387. return ErrNoWebAPISession
  388. }
  389. session.Touch()
  390. return nil
  391. }
  392. // GetAllSessions returns all active sessions (for monitoring/admin).
  393. func (m *WebAPISessionManager) GetAllSessions(ctx context.Context) []*WebAPISession {
  394. m.mu.RLock()
  395. defer m.mu.RUnlock()
  396. sessions := make([]*WebAPISession, 0, len(m.sessions))
  397. for _, session := range m.sessions {
  398. if !session.IsExpired() {
  399. sessions = append(sessions, session)
  400. }
  401. }
  402. return sessions
  403. }
  404. // GetSessionsByScreenName returns all sessions for a given screen name.
  405. func (m *WebAPISessionManager) GetSessionsByScreenName(ctx context.Context, screenName DisplayScreenName) []*WebAPISession {
  406. m.mu.RLock()
  407. defer m.mu.RUnlock()
  408. var sessions []*WebAPISession
  409. identScreenName := screenName.IdentScreenName()
  410. // Check both the byUser map and iterate through all sessions
  411. // since a user might have multiple sessions
  412. for _, session := range m.sessions {
  413. if session.ScreenName.IdentScreenName() == identScreenName {
  414. sessions = append(sessions, session)
  415. }
  416. }
  417. return sessions
  418. }
  419. // cleanupExpiredSessions periodically removes expired sessions.
  420. func (m *WebAPISessionManager) cleanupExpiredSessions() {
  421. for {
  422. select {
  423. case <-m.cleanupTicker.C:
  424. m.mu.Lock()
  425. now := time.Now()
  426. var toRemove []string
  427. for aimsid, session := range m.sessions {
  428. if now.After(session.ExpiresAt) {
  429. toRemove = append(toRemove, aimsid)
  430. }
  431. }
  432. for _, aimsid := range toRemove {
  433. session := m.sessions[aimsid]
  434. delete(m.sessions, aimsid)
  435. delete(m.byUser, session.ScreenName.IdentScreenName())
  436. if session.EventQueue != nil {
  437. session.EventQueue.Close()
  438. }
  439. }
  440. m.mu.Unlock()
  441. case <-m.stopCleanup:
  442. m.cleanupTicker.Stop()
  443. return
  444. }
  445. }
  446. }
  447. // Shutdown stops the session manager and cleans up resources.
  448. func (m *WebAPISessionManager) Shutdown(ctx context.Context) {
  449. close(m.stopCleanup)
  450. m.mu.Lock()
  451. defer m.mu.Unlock()
  452. // CloseSession all event queues
  453. for _, session := range m.sessions {
  454. if session.EventQueue != nil {
  455. session.EventQueue.Close()
  456. }
  457. }
  458. // Clear all sessions
  459. m.sessions = make(map[string]*WebAPISession)
  460. m.byUser = make(map[IdentScreenName]*WebAPISession)
  461. }
  462. // generateSessionID creates a cryptographically secure session ID.
  463. func generateSessionID() (string, error) {
  464. bytes := make([]byte, 32) // 256 bits
  465. if _, err := rand.Read(bytes); err != nil {
  466. return "", err
  467. }
  468. return hex.EncodeToString(bytes), nil
  469. }