package api import ( "encoding/json" "fmt" "net/http" "sort" "strconv" "strings" "time" "browser.local/platform/domain" "browser.local/platform/dto" "browser.local/platform/service" ) const ( defaultLogEventHistoryLimit = 100 maxLogEventHistoryLimit = 10000 logEventHeartbeatInterval = 15 * time.Second ) // serverLogEvents is kept as legacy service plumbing but is not registered as a browser product route. func (h *coreHandlers) serverLogEvents(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { writeMethodNotAllowed(w, http.MethodGet) return } instance, streams, subscription, err := h.openLogEventSubscription(r) if err != nil { writeServiceError(w, err) return } defer subscription.Close() flusher, ok := w.(http.Flusher) if !ok { writeServiceError(w, fmt.Errorf("streaming response unsupported")) return } header := w.Header() header.Set("Content-Type", "text/event-stream") header.Set("Cache-Control", "no-cache, no-transform") header.Set("Connection", "keep-alive") header.Set("X-Accel-Buffering", "no") w.WriteHeader(http.StatusOK) historyLimit := parseLogEventHistoryLimit(r.URL.Query().Get("historyLimit")) for _, stream := range streams { if err := writeSSEJSON(w, "stream", "", dto.LogStreamFromDomain(stream)); err != nil { return } } if historyLimit > 0 { history, err := h.loadLogEventHistory(streams, historyLimit) if err != nil { _ = writeSSEJSON(w, "error", "", map[string]string{"message": err.Error()}) flusher.Flush() return } for _, event := range history { if err := writeSSEJSON(w, "log", logEventID(event), dto.LogStreamEventFromDomain(event)); err != nil { return } } } if err := writeSSEJSON(w, "ready", "", dto.LogStreamEventsReadyResponse{ServerInstanceID: instance.ID, StreamCount: len(streams), ServerTime: time.Now().UTC()}); err != nil { return } flusher.Flush() heartbeat := time.NewTicker(logEventHeartbeatInterval) defer heartbeat.Stop() for { select { case <-r.Context().Done(): return case event, ok := <-subscription.Events: if !ok { return } if err := writeSSEJSON(w, "log", logEventID(event), dto.LogStreamEventFromDomain(event)); err != nil { return } flusher.Flush() case <-heartbeat.C: if _, err := fmt.Fprintf(w, ": heartbeat %s\n\n", time.Now().UTC().Format(time.RFC3339)); err != nil { return } flusher.Flush() } } } func (h *coreHandlers) openLogEventSubscription(r *http.Request) (domain.ServerInstance, []domain.LogStream, service.LogEventSubscription, error) { var instance domain.ServerInstance var streams []domain.LogStream var subscription service.LogEventSubscription var err error if h.enforceAuthorization { sessionID := bearerToken(r) instance, err = h.core.GetServerInstanceForSession(sessionID, r.PathValue("id")) if err != nil { return domain.ServerInstance{}, nil, subscription, err } streams, err = h.core.ListLogStreamsForSession(sessionID, domain.LogStreamFilter{ServerInstanceID: instance.ID}) if err != nil { return domain.ServerInstance{}, nil, subscription, err } subscription, err = h.core.SubscribeLogEventsForSession(sessionID, instance.ID) } else { instance, err = h.core.GetServerInstance(r.PathValue("id")) if err != nil { return domain.ServerInstance{}, nil, subscription, err } streams, err = h.core.ListLogStreams(domain.LogStreamFilter{ServerInstanceID: instance.ID}) if err != nil { return domain.ServerInstance{}, nil, subscription, err } subscription, err = h.core.SubscribeLogEvents(instance.ID) } return instance, streams, subscription, err } func (h *coreHandlers) loadLogEventHistory(streams []domain.LogStream, limit int) ([]domain.LogStreamEvent, error) { history := make([]domain.LogStreamEvent, 0, limit) for _, stream := range streams { afterSeq := uint64(0) if stream.LatestSeq > uint64(limit) { afterSeq = stream.LatestSeq - uint64(limit) } cursor, err := h.core.QueryLogStream(domain.LogStreamCursorQuery{LogStreamID: stream.ID, AfterSeq: afterSeq, Limit: limit}) if err != nil { return nil, err } for _, entry := range cursor.Entries { history = append(history, domain.LogStreamEvent{ServerInstanceID: stream.ServerInstanceID, Stream: stream, Entry: entry, LatestSeq: cursor.LatestSeq}) } } sort.SliceStable(history, func(i, j int) bool { left, right := history[i], history[j] if !left.Entry.Timestamp.Equal(right.Entry.Timestamp) { return left.Entry.Timestamp.Before(right.Entry.Timestamp) } if left.Entry.Seq != right.Entry.Seq { return left.Entry.Seq < right.Entry.Seq } return left.Stream.ID < right.Stream.ID }) if len(history) > limit { history = history[len(history)-limit:] } return history, nil } func parseLogEventHistoryLimit(value string) int { if strings.TrimSpace(value) == "" { return defaultLogEventHistoryLimit } limit, err := strconv.Atoi(value) if err != nil || limit < 0 { return defaultLogEventHistoryLimit } if limit > maxLogEventHistoryLimit { return maxLogEventHistoryLimit } return limit } func writeSSEJSON(w http.ResponseWriter, eventName string, id string, value any) error { payload, err := json.Marshal(value) if err != nil { return err } if id != "" { if _, err := fmt.Fprintf(w, "id: %s\n", sanitizeSSEField(id)); err != nil { return err } } if _, err := fmt.Fprintf(w, "event: %s\n", sanitizeSSEField(eventName)); err != nil { return err } _, err = fmt.Fprintf(w, "data: %s\n\n", payload) return err } func sanitizeSSEField(value string) string { value = strings.ReplaceAll(value, "\r", "") value = strings.ReplaceAll(value, "\n", "") return value } func logEventID(event domain.LogStreamEvent) string { return fmt.Sprintf("%s:%d", sanitizeSSEField(event.Stream.ID), event.Entry.Seq) }