events.go 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196
  1. package handlers
  2. import (
  3. "context"
  4. "encoding/xml"
  5. "fmt"
  6. "log/slog"
  7. "net/http"
  8. "strconv"
  9. "strings"
  10. "time"
  11. "github.com/mk6i/open-oscar-server/server/webapi/types"
  12. "github.com/mk6i/open-oscar-server/state"
  13. )
  14. // EventsHandler handles Web AIM API event fetching endpoints.
  15. type EventsHandler struct {
  16. SessionManager *state.WebAPISessionManager
  17. Logger *slog.Logger
  18. }
  19. // FetchEventsResponse represents the response for fetchEvents endpoint.
  20. type FetchEventsResponse struct {
  21. Response struct {
  22. StatusCode int `json:"statusCode"`
  23. StatusText string `json:"statusText"`
  24. Data FetchEventsData `json:"data"`
  25. } `json:"response"`
  26. }
  27. // FetchEventsData contains the events and metadata.
  28. type FetchEventsData struct {
  29. Events []types.Event `json:"events"`
  30. LastSeqNum uint64 `json:"lastSeqNum"`
  31. TimeToNextFetch int `json:"timeToNextFetch"`
  32. FetchBaseURL string `json:"fetchBaseURL"`
  33. }
  34. // FetchEventsXMLResponse represents the XML response for fetchEvents endpoint.
  35. type FetchEventsXMLResponse struct {
  36. XMLName xml.Name `xml:"response"`
  37. StatusCode int `xml:"statusCode"`
  38. StatusText string `xml:"statusText"`
  39. Data struct {
  40. Events []types.Event `xml:"events>event"`
  41. LastSeqNum uint64 `xml:"lastSeqNum"`
  42. TimeToNextFetch int `xml:"timeToNextFetch"`
  43. FetchBaseURL string `xml:"fetchBaseURL"`
  44. } `xml:"data"`
  45. }
  46. // FetchEvents handles GET /aim/fetchEvents requests with long-polling support.
  47. func (h *EventsHandler) FetchEvents(w http.ResponseWriter, r *http.Request) {
  48. ctx := r.Context()
  49. // Get session ID from parameters
  50. aimsid := r.URL.Query().Get("aimsid")
  51. if aimsid == "" {
  52. h.sendError(w, http.StatusBadRequest, "missing aimsid parameter")
  53. return
  54. }
  55. // Get session
  56. session, err := h.SessionManager.GetSession(r.Context(), aimsid)
  57. if err != nil {
  58. if err == state.ErrNoWebAPISession {
  59. h.sendError(w, http.StatusNotFound, "session not found")
  60. } else if err == state.ErrWebAPISessionExpired {
  61. h.sendError(w, http.StatusGone, "session expired")
  62. } else {
  63. h.sendError(w, http.StatusInternalServerError, "internal server error")
  64. }
  65. return
  66. }
  67. // Touch the session to update last accessed time
  68. h.SessionManager.TouchSession(r.Context(), aimsid)
  69. // Get sequence number parameter
  70. var lastSeqNum uint64
  71. if seqStr := r.URL.Query().Get("seqNum"); seqStr != "" {
  72. if val, err := strconv.ParseUint(seqStr, 10, 64); err == nil {
  73. lastSeqNum = val
  74. }
  75. }
  76. // Get timeout parameter (in seconds, convert to milliseconds)
  77. timeout := time.Duration(session.FetchTimeout) * time.Millisecond
  78. if timeoutStr := r.URL.Query().Get("timeout"); timeoutStr != "" {
  79. if val, err := strconv.Atoi(timeoutStr); err == nil && val > 0 {
  80. timeout = time.Duration(val) * time.Second
  81. }
  82. }
  83. // Limit maximum timeout to 60 seconds
  84. if timeout > 60*time.Second {
  85. timeout = 60 * time.Second
  86. }
  87. // Create a context with timeout for the fetch operation
  88. fetchCtx, cancel := context.WithTimeout(ctx, timeout)
  89. defer cancel()
  90. // Fetch events from the queue (will block until events available or timeout)
  91. events, err := session.EventQueue.Fetch(fetchCtx, lastSeqNum, timeout)
  92. if err != nil {
  93. if err == context.DeadlineExceeded {
  94. // Timeout is normal - return empty events array
  95. events = []types.Event{}
  96. } else {
  97. h.Logger.ErrorContext(ctx, "failed to fetch events", "err", err.Error())
  98. h.sendError(w, http.StatusInternalServerError, "failed to fetch events")
  99. return
  100. }
  101. }
  102. // Determine the last sequence number
  103. var newLastSeqNum uint64 = lastSeqNum
  104. if len(events) > 0 {
  105. newLastSeqNum = events[len(events)-1].SeqNum
  106. }
  107. // Prepare response
  108. resp := FetchEventsResponse{}
  109. resp.Response.StatusCode = 200
  110. resp.Response.StatusText = "OK"
  111. resp.Response.Data.Events = events
  112. resp.Response.Data.LastSeqNum = newLastSeqNum
  113. resp.Response.Data.TimeToNextFetch = session.TimeToNextFetch
  114. // Include fetchBaseURL with updated sequence number for next request
  115. resp.Response.Data.FetchBaseURL = fmt.Sprintf("http://%s/aim/fetchEvents?aimsid=%s&seqNum=%d",
  116. r.Host, aimsid, newLastSeqNum)
  117. // Check response format
  118. format := strings.ToLower(r.URL.Query().Get("f"))
  119. if format == "xml" {
  120. // Send XML response
  121. xmlResp := FetchEventsXMLResponse{}
  122. xmlResp.StatusCode = 200
  123. xmlResp.StatusText = "OK"
  124. xmlResp.Data.Events = events
  125. xmlResp.Data.LastSeqNum = newLastSeqNum
  126. xmlResp.Data.TimeToNextFetch = session.TimeToNextFetch
  127. xmlResp.Data.FetchBaseURL = fmt.Sprintf("http://%s/aim/fetchEvents?aimsid=%s&seqNum=%d",
  128. r.Host, aimsid, newLastSeqNum)
  129. w.Header().Set("Content-Type", "text/xml")
  130. fmt.Fprint(w, `<?xml version="1.0" encoding="UTF-8"?>`)
  131. if err := xml.NewEncoder(w).Encode(xmlResp); err != nil {
  132. h.Logger.Error("failed to encode XML response", "error", err)
  133. }
  134. } else if format == "amf" || format == "amf3" {
  135. // For AMF3, build the response with fields in the correct order
  136. // The working implementation has: response { data {...}, statusCode, statusText, statusDetailCode }
  137. // Convert events to ensure timestamps are float64 for AMF3
  138. convertedEvents := ConvertEventsForAMF3(events)
  139. amfResp := map[string]interface{}{
  140. "response": map[string]interface{}{
  141. // Data comes FIRST (Gromit processes this large object)
  142. "data": map[string]interface{}{
  143. "events": convertedEvents,
  144. "lastSeqNum": newLastSeqNum,
  145. "timeToNextFetch": session.TimeToNextFetch,
  146. "fetchBaseURL": fmt.Sprintf("http://%s/aim/fetchEvents?aimsid=%s&seqNum=%d",
  147. r.Host, aimsid, newLastSeqNum),
  148. },
  149. // Status fields come AFTER data
  150. "statusCode": 200,
  151. "statusText": "OK",
  152. "statusDetailCode": 0,
  153. },
  154. }
  155. // Use SendResponse which will detect AMF format and encode properly
  156. SendResponse(w, r, amfResp, h.Logger)
  157. } else {
  158. // Send JSON/JSONP response with standard structure
  159. SendResponse(w, r, resp, h.Logger)
  160. }
  161. if len(events) > 0 {
  162. h.Logger.DebugContext(ctx, "events fetched",
  163. "aimsid", aimsid,
  164. "count", len(events),
  165. "last_seq", newLastSeqNum,
  166. )
  167. }
  168. }
  169. // sendError is a convenience method that wraps the common SendError function.
  170. func (h *EventsHandler) sendError(w http.ResponseWriter, statusCode int, message string) {
  171. SendError(w, statusCode, message)
  172. }