webapi_session.go 15 KB

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