Current section
Files
Jump to
Current section
Files
src/nimiq_rpc@internal@fiber@src@fiber@backend.erl
-module(nimiq_rpc@internal@fiber@src@fiber@backend).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam").
-export([continue/1, stop/1, build_state/1, wrap/1, stop_on_error/2, handle_text/2, fiber_message/3, handle_binary/2]).
-export_type([direction/0, fiber_builder/0, client_state/0, server_state/0, fiber_state/0, message/0, next/2]).
-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(false).
-type direction() :: server_only_direction |
client_only_direction |
bidirectional_direction.
-type fiber_builder() :: {fiber_builder,
gleam@dict:dict(binary(), fun((gleam@option:option(gleam@dynamic:dynamic_())) -> {ok,
gleam@json:json()} |
{error, nimiq_rpc@internal@fiber@src@fiber@response:error()})),
gleam@dict:dict(binary(), fun((gleam@option:option(gleam@dynamic:dynamic_())) -> nil)),
gleam@option:option(direction())}.
-type client_state() :: {client_state,
gleam@dict:dict(nimiq_rpc@internal@fiber@src@fiber@message:id(), gleam@erlang@process:subject({ok,
gleam@dynamic:dynamic_()} |
{error,
nimiq_rpc@internal@fiber@src@fiber@message:error_data(gleam@dynamic:dynamic_())})),
gleam@dict:dict(gleam@set:set(nimiq_rpc@internal@fiber@src@fiber@message:id()), gleam@erlang@process:subject(gleam@dict:dict(nimiq_rpc@internal@fiber@src@fiber@message:id(), {ok,
gleam@dynamic:dynamic_()} |
{error,
nimiq_rpc@internal@fiber@src@fiber@message:error_data(gleam@dynamic:dynamic_())})))}.
-type server_state() :: {server_state,
gleam@dict:dict(binary(), fun((gleam@option:option(gleam@dynamic:dynamic_())) -> {ok,
gleam@json:json()} |
{error, nimiq_rpc@internal@fiber@src@fiber@response:error()})),
gleam@dict:dict(binary(), fun((gleam@option:option(gleam@dynamic:dynamic_())) -> nil))}.
-opaque fiber_state() :: {client_only, client_state()} |
{server_only, server_state()} |
{bidirectional, client_state(), server_state()}.
-type message() :: {request,
binary(),
gleam@option:option(gleam@json:json()),
nimiq_rpc@internal@fiber@src@fiber@message:id(),
gleam@erlang@process:subject({ok, gleam@dynamic:dynamic_()} |
{error,
nimiq_rpc@internal@fiber@src@fiber@message:error_data(gleam@dynamic:dynamic_())})} |
{notification, binary(), gleam@option:option(gleam@json:json())} |
{batch,
list({binary(),
gleam@option:option(gleam@json:json()),
gleam@option:option(nimiq_rpc@internal@fiber@src@fiber@message:id())}),
gleam@set:set(nimiq_rpc@internal@fiber@src@fiber@message:id()),
gleam@erlang@process:subject(gleam@dict:dict(nimiq_rpc@internal@fiber@src@fiber@message:id(), {ok,
gleam@dynamic:dynamic_()} |
{error,
nimiq_rpc@internal@fiber@src@fiber@message:error_data(gleam@dynamic:dynamic_())}))} |
{remove_waiting, nimiq_rpc@internal@fiber@src@fiber@message:id()} |
{remove_waiting_batch,
gleam@set:set(nimiq_rpc@internal@fiber@src@fiber@message:id())} |
close.
-type next(HTO, HTP) :: {continue,
HTO,
gleam@option:option(gleam@erlang@process:selector(HTP))} |
{stop, gleam@otp@actor:next(HTO, HTP)}.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 93).
?DOC(false).
-spec continue(HUF) -> next(HUF, any()).
continue(State) ->
{continue, State, none}.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 97).
?DOC(false).
-spec stop(gleam@otp@actor:next(HUJ, HUK)) -> next(HUJ, HUK).
stop(Next) ->
{stop, Next}.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 104).
?DOC(false).
-spec build_state(fiber_builder()) -> fiber_state().
build_state(Builder) ->
Direction = gleam@option:lazy_unwrap(
erlang:element(4, Builder),
fun() ->
case gleam@dict:is_empty(erlang:element(2, Builder)) andalso gleam@dict:is_empty(
erlang:element(3, Builder)
) of
true ->
client_only_direction;
false ->
server_only_direction
end
end
),
case Direction of
bidirectional_direction ->
{bidirectional,
{client_state, maps:new(), maps:new()},
{server_state,
erlang:element(2, Builder),
erlang:element(3, Builder)}};
client_only_direction ->
{client_only, {client_state, maps:new(), maps:new()}};
server_only_direction ->
{server_only,
{server_state,
erlang:element(2, Builder),
erlang:element(3, Builder)}}
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 139).
?DOC(false).
-spec wrap(
fun((fun(() -> gleam@erlang@process:selector(message()))) -> {ok, any()} |
{error, HUR})
) -> {ok, gleam@erlang@process:subject(message())} | {error, HUR}.
wrap(Establish) ->
Send_back = gleam@erlang@process:new_subject(),
Bind_selector = fun() ->
Subject = gleam@erlang@process:new_subject(),
gleam@erlang@process:send(Send_back, Subject),
_pipe = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:select(_pipe, Subject)
end,
gleam@result:map(
Establish(Bind_selector),
fun(_) -> _pipe@1 = gleam_erlang_ffi:new_selector(),
_pipe@2 = gleam@erlang@process:select(_pipe@1, Send_back),
gleam_erlang_ffi:select(_pipe@2) end
).
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 160).
?DOC(false).
-spec stop_on_error({ok, any()} | {error, any()}, HVA) -> next(HVA, any()).
stop_on_error(Result, State) ->
case Result of
{error, E} ->
Reason = gleam@string:inspect(E),
gleam_stdlib:print_error(
<<"Fiber closed due to "/utf8, Reason/binary>>
),
stop(gleam@otp@actor:stop_abnormal(Reason));
{ok, _} ->
continue(State)
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 171).
?DOC(false).
-spec add_waiting(
fiber_state(),
nimiq_rpc@internal@fiber@src@fiber@message:id(),
gleam@erlang@process:subject({ok, gleam@dynamic:dynamic_()} |
{error,
nimiq_rpc@internal@fiber@src@fiber@message:error_data(gleam@dynamic:dynamic_())})
) -> fiber_state().
add_waiting(Connection, Id, Reply) ->
case Connection of
{bidirectional, {client_state, Waiting, Waiting_batches}, Server_state} ->
{bidirectional,
{client_state,
begin
_pipe = Waiting,
gleam@dict:insert(_pipe, Id, Reply)
end,
Waiting_batches},
Server_state};
{client_only, {client_state, Waiting@1, Waiting_batches@1}} ->
{client_only,
{client_state,
begin
_pipe@1 = Waiting@1,
gleam@dict:insert(_pipe@1, Id, Reply)
end,
Waiting_batches@1}};
_ ->
Connection
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 194).
?DOC(false).
-spec add_waiting_batch(
fiber_state(),
gleam@set:set(nimiq_rpc@internal@fiber@src@fiber@message:id()),
gleam@erlang@process:subject(gleam@dict:dict(nimiq_rpc@internal@fiber@src@fiber@message:id(), {ok,
gleam@dynamic:dynamic_()} |
{error,
nimiq_rpc@internal@fiber@src@fiber@message:error_data(gleam@dynamic:dynamic_())}))
) -> fiber_state().
add_waiting_batch(Connection, Ids, Reply) ->
case Connection of
{bidirectional, {client_state, Waiting, Waiting_batches}, Server_state} ->
{bidirectional,
{client_state,
Waiting,
begin
_pipe = Waiting_batches,
gleam@dict:insert(_pipe, Ids, Reply)
end},
Server_state};
{client_only, {client_state, Waiting@1, Waiting_batches@1}} ->
{client_only,
{client_state,
Waiting@1,
begin
_pipe@1 = Waiting_batches@1,
gleam@dict:insert(_pipe@1, Ids, Reply)
end}};
_ ->
Connection
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 221).
?DOC(false).
-spec remove_waiting(
fiber_state(),
nimiq_rpc@internal@fiber@src@fiber@message:id()
) -> fiber_state().
remove_waiting(Connection, Id) ->
case Connection of
{bidirectional, {client_state, Waiting, Waiting_batches}, Server_state} ->
{bidirectional,
{client_state,
begin
_pipe = Waiting,
gleam@dict:delete(_pipe, Id)
end,
Waiting_batches},
Server_state};
{client_only, {client_state, Waiting@1, Waiting_batches@1}} ->
{client_only,
{client_state,
begin
_pipe@1 = Waiting@1,
gleam@dict:delete(_pipe@1, Id)
end,
Waiting_batches@1}};
_ ->
Connection
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 240).
?DOC(false).
-spec remove_waiting_batch(
fiber_state(),
gleam@set:set(nimiq_rpc@internal@fiber@src@fiber@message:id())
) -> fiber_state().
remove_waiting_batch(Connection, Ids) ->
case Connection of
{bidirectional, {client_state, Waiting, Waiting_batches}, Server_state} ->
{bidirectional,
{client_state,
Waiting,
begin
_pipe = Waiting_batches,
gleam@dict:delete(_pipe, Ids)
end},
Server_state};
{client_only, {client_state, Waiting@1, Waiting_batches@1}} ->
{client_only,
{client_state,
Waiting@1,
begin
_pipe@1 = Waiting_batches@1,
gleam@dict:delete(_pipe@1, Ids)
end}};
_ ->
Connection
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 429).
?DOC(false).
-spec handle_response(
client_state(),
nimiq_rpc@internal@fiber@src@fiber@message:response(gleam@dynamic:dynamic_())
) -> nil.
handle_response(Client_state, Response) ->
case Response of
{error_response, Error, Id} ->
case begin
_pipe = erlang:element(2, Client_state),
gleam_stdlib:map_get(_pipe, Id)
end of
{error, nil} ->
gleam_stdlib:println_error(
<<<<<<"Received error for id that we were not waiting for (it possibly timed out): "/utf8,
(gleam@string:inspect(Id))/binary>>/binary,
", error: "/utf8>>/binary,
(gleam@string:inspect(Error))/binary>>
);
{ok, Reply_subject} ->
_pipe@1 = Reply_subject,
gleam@erlang@process:send(_pipe@1, {error, Error})
end;
{success_response, Result, Id@1} ->
case begin
_pipe@2 = erlang:element(2, Client_state),
gleam_stdlib:map_get(_pipe@2, Id@1)
end of
{error, nil} ->
gleam_stdlib:println_error(
<<<<<<"Received response for id that we were not waiting for (it possibly timed out): "/utf8,
(gleam@string:inspect(Id@1))/binary>>/binary,
", result: "/utf8>>/binary,
(gleam@string:inspect(Result))/binary>>
);
{ok, Reply_subject@1} ->
_pipe@3 = Reply_subject@1,
gleam@erlang@process:send(_pipe@3, {ok, Result})
end
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 364).
?DOC(false).
-spec handle_request_callback_result(
{ok, gleam@json:json()} |
{error, nimiq_rpc@internal@fiber@src@fiber@response:error()},
nimiq_rpc@internal@fiber@src@fiber@message:id()
) -> nimiq_rpc@internal@fiber@src@fiber@message:response(gleam@json:json()).
handle_request_callback_result(Result, Id) ->
case Result of
{error, invalid_params} ->
_pipe = none,
_pipe@1 = {error_data, _pipe, -32602, <<"Invalid params"/utf8>>},
{error_response, _pipe@1, Id};
{error, internal_error} ->
_pipe@2 = none,
_pipe@3 = {error_data, _pipe@2, -32603, <<"Internal error"/utf8>>},
{error_response, _pipe@3, Id};
{error, {custom_error, Error}} ->
_pipe@4 = Error,
{error_response, _pipe@4, Id};
{ok, Result@1} ->
_pipe@5 = Result@1,
{success_response, _pipe@5, Id}
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 390).
?DOC(false).
-spec process_request(
server_state(),
nimiq_rpc@internal@fiber@src@fiber@message:request(gleam@dynamic:dynamic_())
) -> {ok,
nimiq_rpc@internal@fiber@src@fiber@message:response(gleam@json:json())} |
{error, nil}.
process_request(Server_state, Request) ->
case Request of
{notification, Params, Method} ->
case begin
_pipe = erlang:element(3, Server_state),
gleam_stdlib:map_get(_pipe, Method)
end of
{error, nil} ->
gleam_stdlib:println_error(
<<<<<<"Received notification we don't have a handler for: "/utf8,
Method/binary>>/binary,
", params: "/utf8>>/binary,
(gleam@string:inspect(Params))/binary>>
),
{error, nil};
{ok, Callback} ->
Callback(Params),
{error, nil}
end;
{request, Params@1, Method@1, Id} ->
case begin
_pipe@1 = erlang:element(2, Server_state),
gleam_stdlib:map_get(_pipe@1, Method@1)
end of
{error, nil} ->
_pipe@2 = {some, gleam@json:string(Method@1)},
_pipe@3 = {error_data,
_pipe@2,
-32601,
<<"Method not found"/utf8>>},
_pipe@4 = {error_response, _pipe@3, Id},
{ok, _pipe@4};
{ok, Callback@1} ->
_pipe@5 = Callback@1(Params@1),
_pipe@6 = handle_request_callback_result(_pipe@5, Id),
{ok, _pipe@6}
end
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 485).
?DOC(false).
-spec handle_batch_response(
client_state(),
list(nimiq_rpc@internal@fiber@src@fiber@message:response(gleam@dynamic:dynamic_()))
) -> nil.
handle_batch_response(Client_state, Batch) ->
Ids = begin
_pipe = Batch,
_pipe@1 = gleam@list:map(_pipe, fun(Response) -> case Response of
{error_response, _, Id} ->
Id;
{success_response, _, Id@1} ->
Id@1
end end),
gleam@set:from_list(_pipe@1)
end,
case begin
_pipe@2 = erlang:element(3, Client_state),
gleam_stdlib:map_get(_pipe@2, Ids)
end of
{error, nil} ->
gleam_stdlib:println_error(
<<"Received batch response for an id set that we were not waiting for (it possibly timed out): "/utf8,
(gleam@string:inspect(Ids))/binary>>
);
{ok, Reply_subject} ->
_pipe@3 = Batch,
_pipe@4 = gleam@list:map(
_pipe@3,
fun(Response@1) -> case Response@1 of
{error_response, Error, Id@2} ->
{Id@2, {error, Error}};
{success_response, Result, Id@3} ->
{Id@3, {ok, Result}}
end end
),
_pipe@5 = maps:from_list(_pipe@4),
gleam@erlang@process:send(Reply_subject, _pipe@5)
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 467).
?DOC(false).
-spec handle_batch_request(
server_state(),
list(nimiq_rpc@internal@fiber@src@fiber@message:request(gleam@dynamic:dynamic_()))
) -> {ok, nimiq_rpc@internal@fiber@src@fiber@message:message(gleam@json:json())} |
{error, nil}.
handle_batch_request(Server_state, Batch) ->
Responses = begin
_pipe = Batch,
_pipe@1 = gleam@list:map(
_pipe,
fun(_capture) -> process_request(Server_state, _capture) end
),
gleam@result:values(_pipe@1)
end,
case Responses of
[] ->
{error, nil};
Responses@1 ->
_pipe@2 = Responses@1,
_pipe@3 = {batch_response_message, _pipe@2},
{ok, _pipe@3}
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 521).
?DOC(false).
-spec handle_message(
fiber_state(),
nimiq_rpc@internal@fiber@src@fiber@message:message(gleam@dynamic:dynamic_())
) -> {ok, nimiq_rpc@internal@fiber@src@fiber@message:message(gleam@json:json())} |
{error, nil}.
handle_message(State, Fiber_message) ->
case Fiber_message of
{batch_request_message, Batch} ->
case State of
{server_only, Server_state} ->
handle_batch_request(Server_state, Batch);
{bidirectional, _, Server_state} ->
handle_batch_request(Server_state, Batch);
_ ->
{error, nil}
end;
{batch_response_message, Batch@1} ->
case State of
{client_only, Client_state} ->
{error, handle_batch_response(Client_state, Batch@1)};
{bidirectional, Client_state, _} ->
{error, handle_batch_response(Client_state, Batch@1)};
_ ->
{error, nil}
end;
{request_message, Request} ->
case State of
{server_only, Server_state@1} ->
case process_request(Server_state@1, Request) of
{error, nil} ->
{error, nil};
{ok, Response} ->
_pipe = Response,
_pipe@1 = {response_message, _pipe},
{ok, _pipe@1}
end;
{bidirectional, _, Server_state@1} ->
case process_request(Server_state@1, Request) of
{error, nil} ->
{error, nil};
{ok, Response} ->
_pipe = Response,
_pipe@1 = {response_message, _pipe},
{ok, _pipe@1}
end;
_ ->
{error, nil}
end;
{response_message, Response@1} ->
case State of
{client_only, Client_state@1} ->
{error, handle_response(Client_state@1, Response@1)};
{bidirectional, Client_state@1, _} ->
{error, handle_response(Client_state@1, Response@1)};
_ ->
{error, nil}
end;
{error_message, Error} ->
gleam_stdlib:println_error(
<<"Received error without id (usually indicates we sent malformed data): "/utf8,
(gleam@string:inspect(Error))/binary>>
),
{error, nil}
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 262).
?DOC(false).
-spec handle_text(fiber_state(), binary()) -> {ok,
nimiq_rpc@internal@fiber@src@fiber@message:message(gleam@json:json())} |
{error, nil}.
handle_text(State, Text) ->
case nimiq_rpc@internal@fiber@src@fiber@message:decode(Text) of
{error, Error} ->
_pipe = case Error of
{unable_to_decode, _} ->
{error_data, none, -32600, <<"Invalid Request"/utf8>>};
{unexpected_byte, Byte} ->
{error_data,
{some,
gleam@json:string(
<<<<"Unexpected Byte: \""/utf8, Byte/binary>>/binary,
"\""/utf8>>
)},
-32700,
<<"Parse error"/utf8>>};
unexpected_end_of_input ->
{error_data,
{some,
gleam@json:string(
<<"Unexpected End of Input"/utf8>>
)},
-32700,
<<"Parse error"/utf8>>};
{unexpected_sequence, Sequence} ->
{error_data,
{some,
gleam@json:string(
<<<<"Unexpected Sequence: \""/utf8,
Sequence/binary>>/binary,
"\""/utf8>>
)},
-32700,
<<"Parse error"/utf8>>}
end,
_pipe@1 = {error_message, _pipe},
{ok, _pipe@1};
{ok, Message} ->
handle_message(State, Message)
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 304).
?DOC(false).
-spec fiber_message(
fiber_state(),
message(),
fun((nimiq_rpc@internal@fiber@src@fiber@message:message(gleam@json:json())) -> {ok,
any()} |
{error, any()})
) -> next(fiber_state(), any()).
fiber_message(State, Message, Send) ->
case Message of
{request, Method, Params, Id, Reply_subject} ->
_pipe = {request, Params, Method, Id},
_pipe@1 = {request_message, _pipe},
_pipe@2 = Send(_pipe@1),
stop_on_error(
_pipe@2,
begin
_pipe@3 = State,
add_waiting(_pipe@3, Id, Reply_subject)
end
);
{notification, Method@1, Params@1} ->
_pipe@4 = {notification, Params@1, Method@1},
_pipe@5 = {request_message, _pipe@4},
_pipe@6 = Send(_pipe@5),
stop_on_error(_pipe@6, State);
{batch, Batch, Ids, Reply_subject@1} ->
_pipe@7 = Batch,
_pipe@8 = gleam@list:map(
_pipe@7,
fun(Request) ->
{Method@2, Params@2, Id@1} = Request,
case Id@1 of
none ->
{notification, Params@2, Method@2};
{some, Id@2} ->
{request, Params@2, Method@2, Id@2}
end
end
),
_pipe@9 = {batch_request_message, _pipe@8},
_pipe@10 = Send(_pipe@9),
stop_on_error(
_pipe@10,
begin
_pipe@11 = State,
add_waiting_batch(_pipe@11, Ids, Reply_subject@1)
end
);
{remove_waiting, Id@3} ->
continue(
begin
_pipe@12 = State,
remove_waiting(_pipe@12, Id@3)
end
);
{remove_waiting_batch, Ids@1} ->
continue(
begin
_pipe@13 = State,
remove_waiting_batch(_pipe@13, Ids@1)
end
);
close ->
stop(gleam@otp@actor:stop())
end.
-file("src/nimiq_rpc/internal/fiber/src/fiber/backend.gleam", 341).
?DOC(false).
-spec handle_binary(fiber_state(), bitstring()) -> {ok,
nimiq_rpc@internal@fiber@src@fiber@message:message(gleam@json:json())} |
{error, nil}.
handle_binary(State, _) ->
case State of
{bidirectional, _, _} ->
_pipe = {error_data,
{some,
gleam@json:string(<<"binary frames are unsupported"/utf8>>)},
-32700,
<<"Parse error"/utf8>>},
_pipe@1 = {error_message, _pipe},
{ok, _pipe@1};
{server_only, _} ->
_pipe = {error_data,
{some,
gleam@json:string(<<"binary frames are unsupported"/utf8>>)},
-32700,
<<"Parse error"/utf8>>},
_pipe@1 = {error_message, _pipe},
{ok, _pipe@1};
{client_only, _} ->
gleam_stdlib:println_error(
<<"Received binary data, which is unsupported by this backend"/utf8>>
),
{error, nil}
end.