package convert import ( "bufio" "encoding/json" "fmt" "io" "net/http" "strings" ) // SSEWriter writes Server-Sent Events type SSEWriter struct { writer io.Writer flusher http.Flusher } // NewSSEWriter creates a new SSE writer func NewSSEWriter(w http.ResponseWriter) *SSEWriter { flusher, _ := w.(http.Flusher) return &SSEWriter{ writer: w, flusher: flusher, } } // WriteEvent writes a single SSE event func (w *SSEWriter) WriteEvent(event string, data interface{}) error { var dataStr string switch v := data.(type) { case string: dataStr = v default: b, err := json.Marshal(v) if err != nil { return err } dataStr = string(b) } _, err := fmt.Fprintf(w.writer, "event: %s\ndata: %s\n\n", event, dataStr) if err != nil { return err } if w.flusher != nil { w.flusher.Flush() } return nil } // WriteChunk writes a streaming chunk in SSE format func (w *SSEWriter) WriteChunk(chunk interface{}) error { b, err := json.Marshal(chunk) if err != nil { return err } _, err = fmt.Fprintf(w.writer, "data: %s\n\n", string(b)) if err != nil { return err } if w.flusher != nil { w.flusher.Flush() } return nil } // WriteDone writes the [DONE] marker func (w *SSEWriter) WriteDone() error { _, err := fmt.Fprintf(w.writer, "data: [DONE]\n\n") if err != nil { return err } if w.flusher != nil { w.flusher.Flush() } return nil } // SSEParser parses Server-Sent Events from a reader type SSEParser struct { reader *bufio.Reader } // NewSSEParser creates a new SSE parser func NewSSEParser(r io.Reader) *SSEParser { return &SSEParser{ reader: bufio.NewReader(r), } } // SSEEvent represents a parsed SSE event type SSEEvent struct { Event string Data string } // ReadEvent reads the next SSE event func (p *SSEParser) ReadEvent() (*SSEEvent, error) { event := &SSEEvent{} for { line, err := p.reader.ReadString('\n') if err != nil { return nil, err } line = strings.TrimRight(line, "\r\n") if line == "" { // Empty line means end of event if event.Data != "" || event.Event != "" { return event, nil } continue } if strings.HasPrefix(line, "event:") { event.Event = strings.TrimSpace(line[6:]) } else if strings.HasPrefix(line, "data:") { data := strings.TrimSpace(line[5:]) if event.Data != "" { event.Data += "\n" + data } else { event.Data = data } } // Ignore comments (lines starting with :) and unknown fields } } // ParseChatStreamChunk parses an OpenAI Chat Completions streaming chunk func ParseChatStreamChunk(data string) (*ChatCompletionStreamChunk, error) { if data == "[DONE]" { return nil, io.EOF } var chunk ChatCompletionStreamChunk err := json.Unmarshal([]byte(data), &chunk) if err != nil { return nil, err } return &chunk, nil } // ParseMessagesStreamEvent parses an Anthropic Messages streaming event func ParseMessagesStreamEvent(data string) (*AnthropicStreamEvent, error) { var event AnthropicStreamEvent err := json.Unmarshal([]byte(data), &event) if err != nil { return nil, err } return &event, nil } // ParseResponsesStreamChunk parses an OpenAI Responses API streaming chunk func ParseResponsesStreamChunk(data string) (*ResponsesStreamEvent, error) { if data == "[DONE]" { return nil, io.EOF } var event ResponsesStreamEvent err := json.Unmarshal([]byte(data), &event) if err != nil { return nil, err } return &event, nil }