package api import ( "encoding/json" "fmt" "net/http" "strconv" "strings" "time" "browser.local/platform/domain" "browser.local/platform/dto" "browser.local/platform/service" ) const ( defaultLogEventHistoryLimit = 100 maxLogEventHistoryLimit = 500 logEventHeartbeatInterval = 15 * time.Second ) // serverLogEvents godoc // @Summary Stream server log events // @Description Streams safe server log entries over Server-Sent Events. Durable cursor query remains available for history and reconnect repair. // @Tags logs // @Produce text/event-stream // @Param id path string true "Server instance ID" // @Param historyLimit query int false "Recent entries per stream to replay before live events" // @Success 200 {string} string "event-stream" // @Failure 401 {object} dto.ErrorResponse // @Failure 403 {object} dto.ErrorResponse // @Failure 404 {object} dto.ErrorResponse // @Failure 405 {object} dto.ErrorResponse // @Router /api/v1/server-instances/{id}/logs/events [get] 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 { if err := h.writeLogEventHistory(w, stream, historyLimit); err != nil { _ = writeSSEJSON(w, "error", "", map[string]string{"message": err.Error()}) flusher.Flush() 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) writeLogEventHistory(w http.ResponseWriter, stream domain.LogStream, limit int) error { 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 err } for _, entry := range cursor.Entries { event := domain.LogStreamEvent{ServerInstanceID: stream.ServerInstanceID, Stream: stream, Entry: entry, LatestSeq: cursor.LatestSeq} if err := writeSSEJSON(w, "log", logEventID(event), dto.LogStreamEventFromDomain(event)); err != nil { return err } } return 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) }