Packages

A Gleam library for interacting with the Corrosion API

Current section

Files

Jump to
corrosion src corrosion@subscription.erl
Raw

src/corrosion@subscription.erl

-module(corrosion@subscription).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/corrosion/subscription.gleam").
-export([subscribe/4, unsubscribe/1]).
-export_type([change_type/0, event/1, message/0, internal_message/0, state/1]).
-type change_type() :: insert | update | delete.
-type event(ITN) :: {row, integer(), ITN} |
{change, change_type(), integer(), ITN, integer()} |
{end_of_query, float(), gleam@option:option(integer())} |
{query_error, binary()} |
{decode_error, list(gleam@dynamic@decode:decode_error())} |
closed.
-type message() :: shutdown.
-type internal_message() :: {control_message, message()} |
{stream_event, httpp@jsonl:jsonl_event(corrosion@query_event:query_event())}.
-type state(ITO) :: {state,
gleam@dynamic@decode:decoder(ITO),
gleam@erlang@process:subject(event(ITO)),
list(gleam@dynamic:dynamic_()),
gleam@erlang@process:subject(httpp@jsonl:jsonl_manager_message())}.
-file("src/corrosion/subscription.gleam", 23).
-spec query_event_change_type_to_event_change_type(
corrosion@query_event:change_type()
) -> change_type().
query_event_change_type_to_event_change_type(In) ->
case In of
delete ->
delete;
insert ->
insert;
update ->
update
end.
-file("src/corrosion/subscription.gleam", 60).
-spec initialize(
gleam@erlang@process:subject(internal_message()),
gleam@uri:uri(),
corrosion@statement:statement(),
gleam@dynamic@decode:decoder(ITQ),
gleam@erlang@process:subject(event(ITQ))
) -> {ok,
gleam@otp@actor:initialised(state(ITQ), internal_message(), gleam@erlang@process:subject(message()))} |
{error, binary()}.
initialize(Self, Corro_uri, Statement, Row_decoder, Recv) ->
Uri = {uri,
erlang:element(2, Corro_uri),
erlang:element(3, Corro_uri),
erlang:element(4, Corro_uri),
erlang:element(5, Corro_uri),
<<"/v1/subscriptions"/utf8>>,
erlang:element(7, Corro_uri),
erlang:element(8, Corro_uri)},
Body = begin
_pipe = corrosion@internal@util:statement_to_json(Statement),
_pipe@1 = gleam_json_ffi:json_to_iodata(_pipe),
gleam_stdlib:wrap_list(_pipe@1)
end,
Base_request@1 = case gleam@http@request:from_uri(Uri) of
{ok, Base_request} -> Base_request;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"corrosion/subscription"/utf8>>,
function => <<"initialize"/utf8>>,
line => 76,
value => _assert_fail,
start => 1732,
'end' => 1783,
pattern_start => 1743,
pattern_end => 1759})
end,
Request = begin
_pipe@2 = Base_request@1,
_pipe@3 = gleam@http@request:set_header(
_pipe@2,
<<"content-type"/utf8>>,
<<"application/json"/utf8>>
),
gleam@http@request:set_body(_pipe@3, Body)
end,
Jsonl_event_subject = gleam@erlang@process:new_subject(),
Req_mgr_result = begin
_pipe@4 = httpp@jsonl:json_lines_stream(
Request,
1000,
corrosion@query_event:decoder(),
Jsonl_event_subject
),
gleam@result:replace_error(
_pipe@4,
<<"Could not start jsonl stream."/utf8>>
)
end,
gleam@result:'try'(
Req_mgr_result,
fun(_use0) ->
{_, Req_mgr} = _use0,
Control_message_subject = gleam@erlang@process:new_subject(),
Initial_state = {state, Row_decoder, Recv, [], Req_mgr},
Selector = begin
_pipe@5 = gleam_erlang_ffi:new_selector(),
_pipe@6 = gleam@erlang@process:select(_pipe@5, Self),
_pipe@7 = gleam@erlang@process:select_map(
_pipe@6,
Jsonl_event_subject,
fun(Field@0) -> {stream_event, Field@0} end
),
gleam@erlang@process:select_map(
_pipe@7,
Control_message_subject,
fun(Field@0) -> {control_message, Field@0} end
)
end,
_pipe@8 = gleam@otp@actor:initialised(Initial_state),
_pipe@9 = gleam@otp@actor:selecting(_pipe@8, Selector),
_pipe@10 = gleam@otp@actor:returning(
_pipe@9,
Control_message_subject
),
{ok, _pipe@10}
end
).
-file("src/corrosion/subscription.gleam", 198).
-spec decode_data(
list(gleam@dynamic:dynamic_()),
list(gleam@dynamic:dynamic_()),
gleam@dynamic@decode:decoder(IUY)
) -> {ok, IUY} | {error, event(IUY)}.
decode_data(Columns, Values, Decoder) ->
Zipped = begin
_pipe = gleam@list:strict_zip(Columns, Values),
_pipe@1 = gleam@result:map(_pipe, fun gleam@dynamic:properties/1),
gleam@result:replace_error(
_pipe@1,
{query_error,
<<"Columns array is a different length from values."/utf8>>}
)
end,
gleam@result:'try'(
Zipped,
fun(Dynamic) ->
Decode_result = begin
_pipe@2 = gleam@dynamic@decode:run(Dynamic, Decoder),
gleam@result:try_recover(
_pipe@2,
fun(Decode_errors) ->
case gleam@dynamic@decode:run(
gleam_stdlib:identity(Values),
Decoder
) of
{ok, Data} ->
{ok, Data};
{error, _} ->
{error, {decode_error, Decode_errors}}
end
end
)
end,
gleam@result:'try'(Decode_result, fun(Result) -> {ok, Result} end)
end
).
-file("src/corrosion/subscription.gleam", 145).
-spec handle_query_event(state(IUR), corrosion@query_event:query_event()) -> gleam@otp@actor:next(state(IUR), internal_message()).
handle_query_event(State, Event) ->
case Event of
{columns, Column_names} ->
gleam@otp@actor:continue(
{state,
erlang:element(2, State),
erlang:element(3, State),
gleam@list:map(Column_names, fun gleam_stdlib:identity/1),
erlang:element(5, State)}
);
{e_o_q, Time, Change_id} ->
gleam@erlang@process:send(
erlang:element(3, State),
{end_of_query, Time, Change_id}
),
gleam@otp@actor:continue(State);
{query_error, Message} ->
gleam@erlang@process:send(
erlang:element(3, State),
{query_error, Message}
),
gleam@otp@actor:continue(State);
{row, Rowid, Values} ->
case decode_data(
erlang:element(4, State),
Values,
erlang:element(2, State)
) of
{ok, Data} ->
gleam@erlang@process:send(
erlang:element(3, State),
{row, Rowid, Data}
),
gleam@otp@actor:continue(State);
{error, Event@1} ->
gleam@erlang@process:send(erlang:element(3, State), Event@1),
gleam@otp@actor:continue(State)
end;
{change, Change_type, Row_id, Values@1, Change_id@1} ->
case decode_data(
erlang:element(4, State),
Values@1,
erlang:element(2, State)
) of
{ok, Data@1} ->
gleam@erlang@process:send(
erlang:element(3, State),
{change,
query_event_change_type_to_event_change_type(
Change_type
),
Row_id,
Data@1,
Change_id@1}
),
gleam@otp@actor:continue(State);
{error, Evt} ->
gleam@erlang@process:send(erlang:element(3, State), Evt),
gleam@otp@actor:continue(State)
end
end.
-file("src/corrosion/subscription.gleam", 226).
-spec handle_shutdown(state(IVD)) -> gleam@otp@actor:next(state(IVD), internal_message()).
handle_shutdown(State) ->
case gleam@erlang@process:subject_owner(erlang:element(3, State)) of
{ok, Pid} ->
gleam@erlang@process:unlink(Pid);
_ ->
nil
end,
gleam@erlang@process:send(erlang:element(5, State), shutdown),
gleam@otp@actor:stop().
-file("src/corrosion/subscription.gleam", 123).
-spec handle_control_message(state(IUG), message()) -> gleam@otp@actor:next(state(IUG), internal_message()).
handle_control_message(State, Message) ->
case Message of
shutdown ->
handle_shutdown(State)
end.
-file("src/corrosion/subscription.gleam", 132).
-spec handle_jsonl_event(
state(IUL),
httpp@jsonl:jsonl_event(corrosion@query_event:query_event())
) -> gleam@otp@actor:next(state(IUL), internal_message()).
handle_jsonl_event(State, Event) ->
case Event of
closed ->
gleam@erlang@process:send(erlang:element(3, State), closed),
handle_shutdown(State);
{line, Evt} ->
handle_query_event(State, Evt)
end.
-file("src/corrosion/subscription.gleam", 113).
-spec on_message(state(IUB), internal_message()) -> gleam@otp@actor:next(state(IUB), internal_message()).
on_message(State, Message) ->
case Message of
{control_message, Msg} ->
handle_control_message(State, Msg);
{stream_event, Evt} ->
handle_jsonl_event(State, Evt)
end.
-file("src/corrosion/subscription.gleam", 239).
-spec subscribe(
gleam@uri:uri(),
corrosion@statement:statement(),
gleam@dynamic@decode:decoder(IVI),
gleam@erlang@process:subject(event(IVI))
) -> {ok, gleam@erlang@process:subject(message())} | {error, nil}.
subscribe(Corro_uri, Statement, Row_decoder, Recv) ->
Actor = begin
_pipe = gleam@otp@actor:new_with_initialiser(
1000,
fun(_capture) ->
initialize(_capture, Corro_uri, Statement, Row_decoder, Recv)
end
),
_pipe@1 = gleam@otp@actor:on_message(_pipe, fun on_message/2),
gleam@otp@actor:start(_pipe@1)
end,
case Actor of
{ok, Started} ->
case erlang:is_process_alive(erlang:element(2, Started)) of
true -> nil;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"corrosion/subscription"/utf8>>,
function => <<"subscribe"/utf8>>,
line => 258,
value => _assert_fail,
start => 6395,
'end' => 6442,
pattern_start => 6406,
pattern_end => 6410})
end,
gleam_erlang_ffi:link(erlang:element(2, Started)),
{ok, erlang:element(3, Started)};
_ ->
{error, nil}
end.
-file("src/corrosion/subscription.gleam", 266).
-spec unsubscribe(gleam@erlang@process:subject(message())) -> nil.
unsubscribe(Subscription) ->
gleam@erlang@process:send(Subscription, shutdown).