-module(anthropic@streaming@sse). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]). -define(FILEPATH, "src/anthropic/streaming/sse.gleam"). -export([new_parser_state/0, reset_event_state/1, is_keepalive/1, get_event_type/1, get_data/1, sse_event/2, sse_event_full/4, parse_line/2, parse_event_lines/1, parse_event/1, parse_chunk/2, flush/1]). -export_type([sse_event/0, sse_parser_state/0, sse_parse_result/0, sse_error/0]). -if(?OTP_RELEASE >= 27). -define(MODULEDOC(Str), -moduledoc(Str)). -define(DOC(Str), -doc(Str)). -else. -define(MODULEDOC(Str), -compile([])). -define(DOC(Str), -compile([])). -endif. ?MODULEDOC( " Server-Sent Events (SSE) parser for Anthropic streaming responses\n" "\n" " This module provides parsing utilities for SSE format used by the\n" " Anthropic streaming API. It handles multi-line data fields, event types,\n" " keepalives, and malformed events gracefully.\n" ). -type sse_event() :: {sse_event, gleam@option:option(binary()), gleam@option:option(binary()), gleam@option:option(binary()), gleam@option:option(integer())}. -type sse_parser_state() :: {sse_parser_state, gleam@option:option(binary()), list(binary()), gleam@option:option(binary()), gleam@option:option(integer()), binary()}. -type sse_parse_result() :: {sse_parse_result, list(sse_event()), sse_parser_state()}. -type sse_error() :: {invalid_field, binary()} | empty_event | {invalid_retry, binary()}. -file("src/anthropic/streaming/sse.gleam", 71). ?DOC(" Create a new parser state\n"). -spec new_parser_state() -> sse_parser_state(). new_parser_state() -> {sse_parser_state, none, [], none, none, <<""/utf8>>}. -file("src/anthropic/streaming/sse.gleam", 82). ?DOC(" Reset parser state for next event (keeps buffer)\n"). -spec reset_event_state(sse_parser_state()) -> sse_parser_state(). reset_event_state(State) -> {sse_parser_state, none, [], none, none, erlang:element(6, State)}. -file("src/anthropic/streaming/sse.gleam", 178). ?DOC(" Build an SseEvent from the current parser state\n"). -spec build_event(sse_parser_state()) -> {ok, sse_event()} | {error, sse_error()}. build_event(State) -> Data = case erlang:element(3, State) of [] -> none; Lines -> {some, gleam@string:join(Lines, <<"\n"/utf8>>)} end, case {Data, erlang:element(2, State)} of {none, none} -> {error, empty_event}; {_, _} -> {ok, {sse_event, erlang:element(2, State), Data, erlang:element(4, State), erlang:element(5, State)}} end. -file("src/anthropic/streaming/sse.gleam", 258). ?DOC(" Check if an SSE event is a keepalive/ping\n"). -spec is_keepalive(sse_event()) -> boolean(). is_keepalive(Event) -> case {erlang:element(2, Event), erlang:element(3, Event)} of {none, none} -> true; {{some, <<"ping"/utf8>>}, _} -> true; {_, _} -> false end. -file("src/anthropic/streaming/sse.gleam", 267). ?DOC(" Get the event type as a string, defaulting to \"message\" if not specified\n"). -spec get_event_type(sse_event()) -> binary(). get_event_type(Event) -> gleam@option:unwrap(erlang:element(2, Event), <<"message"/utf8>>). -file("src/anthropic/streaming/sse.gleam", 272). ?DOC(" Get the data as a string, defaulting to empty string\n"). -spec get_data(sse_event()) -> binary(). get_data(Event) -> gleam@option:unwrap(erlang:element(3, Event), <<""/utf8>>). -file("src/anthropic/streaming/sse.gleam", 277). ?DOC(" Create an SSE event (primarily for testing)\n"). -spec sse_event(gleam@option:option(binary()), gleam@option:option(binary())) -> sse_event(). sse_event(Event_type, Data) -> {sse_event, Event_type, Data, none, none}. -file("src/anthropic/streaming/sse.gleam", 282). ?DOC(" Create an SSE event with all fields\n"). -spec sse_event_full( gleam@option:option(binary()), gleam@option:option(binary()), gleam@option:option(binary()), gleam@option:option(integer()) ) -> sse_event(). sse_event_full(Event_type, Data, Id, Retry) -> {sse_event, Event_type, Data, Id, Retry}. -file("src/anthropic/streaming/sse.gleam", 296). ?DOC(" Parse an integer from a string\n"). -spec parse_int(binary()) -> {ok, integer()} | {error, nil}. parse_int(Str) -> case gleam_stdlib:parse_int(Str) of {ok, N} -> {ok, N}; {error, _} -> {error, nil} end. -file("src/anthropic/streaming/sse.gleam", 154). ?DOC(" Apply a parsed field to the parser state\n"). -spec apply_field(sse_parser_state(), binary(), binary()) -> sse_parser_state(). apply_field(State, Field, Value) -> case Field of <<"event"/utf8>> -> {sse_parser_state, {some, Value}, erlang:element(3, State), erlang:element(4, State), erlang:element(5, State), erlang:element(6, State)}; <<"data"/utf8>> -> {sse_parser_state, erlang:element(2, State), lists:append(erlang:element(3, State), [Value]), erlang:element(4, State), erlang:element(5, State), erlang:element(6, State)}; <<"id"/utf8>> -> {sse_parser_state, erlang:element(2, State), erlang:element(3, State), {some, Value}, erlang:element(5, State), erlang:element(6, State)}; <<"retry"/utf8>> -> case parse_int(Value) of {ok, N} -> {sse_parser_state, erlang:element(2, State), erlang:element(3, State), erlang:element(4, State), {some, N}, erlang:element(6, State)}; {error, _} -> State end; _ -> State end. -file("src/anthropic/streaming/sse.gleam", 129). ?DOC(" Parse a field line (field: value format)\n"). -spec parse_field_line(sse_parser_state(), binary()) -> sse_parser_state(). parse_field_line(State, Line) -> case gleam_stdlib:string_starts_with(Line, <<":"/utf8>>) of true -> State; false -> case gleam@string:split_once(Line, <<":"/utf8>>) of {ok, {Field, Value}} -> Trimmed_value = case gleam_stdlib:string_starts_with( Value, <<" "/utf8>> ) of true -> gleam@string:drop_start(Value, 1); false -> Value end, apply_field(State, Field, Trimmed_value); {error, _} -> apply_field(State, Line, <<""/utf8>>) end end. -file("src/anthropic/streaming/sse.gleam", 120). ?DOC(" Parse a single line and update parser state\n"). -spec parse_line(sse_parser_state(), binary()) -> sse_parser_state(). parse_line(State, Line) -> case gleam@string:trim(Line) of <<""/utf8>> -> State; _ -> parse_field_line(State, Line) end. -file("src/anthropic/streaming/sse.gleam", 103). ?DOC(" Parse SSE event from a list of lines\n"). -spec parse_event_lines(list(binary())) -> {ok, sse_event()} | {error, sse_error()}. parse_event_lines(Lines) -> Initial_state = {sse_parser_state, none, [], none, none, <<""/utf8>>}, Final_state = gleam@list:fold( Lines, Initial_state, fun(State, Line) -> parse_line(State, Line) end ), build_event(Final_state). -file("src/anthropic/streaming/sse.gleam", 97). ?DOC(" Parse a single SSE event from a complete event block (lines separated by \\n)\n"). -spec parse_event(binary()) -> {ok, sse_event()} | {error, sse_error()}. parse_event(Event_block) -> Lines = gleam@string:split(Event_block, <<"\n"/utf8>>), parse_event_lines(Lines). -file("src/anthropic/streaming/sse.gleam", 206). ?DOC( " Parse a chunk of SSE data, returning parsed events and updated state\n" "\n" " This function handles partial events that span multiple chunks by\n" " maintaining a buffer in the parser state.\n" ). -spec parse_chunk(sse_parser_state(), binary()) -> sse_parse_result(). parse_chunk(State, Chunk) -> Full_data = <<(erlang:element(6, State))/binary, Chunk/binary>>, Parts = gleam@string:split(Full_data, <<"\n\n"/utf8>>), {Events, Remaining} = case Parts of [] -> {[], <<""/utf8>>}; [Single] -> {[], Single}; _ -> Complete_parts = gleam@list:take(Parts, erlang:length(Parts) - 1), Last_part = begin _pipe = gleam@list:last(Parts), gleam@result:unwrap(_pipe, <<""/utf8>>) end, Parsed_events = begin _pipe@1 = Complete_parts, gleam@list:filter_map( _pipe@1, fun(Part) -> case gleam@string:trim(Part) of <<""/utf8>> -> {error, nil}; Trimmed -> case parse_event(Trimmed) of {ok, Event} -> {ok, Event}; {error, _} -> {error, nil} end end end ) end, {Parsed_events, Last_part} end, {sse_parse_result, Events, begin _record = new_parser_state(), {sse_parser_state, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), erlang:element(5, _record), Remaining} end}. -file("src/anthropic/streaming/sse.gleam", 246). ?DOC(" Flush any remaining data in the buffer as a final event\n"). -spec flush(sse_parser_state()) -> {ok, sse_event()} | {error, sse_error()}. flush(State) -> case gleam@string:trim(erlang:element(6, State)) of <<""/utf8>> -> {error, empty_event}; Data -> parse_event(Data) end.