Current section

Files

Jump to
telega src telega@flow@compose.erl
Raw

src/telega@flow@compose.erl

-module(telega@flow@compose).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/telega/flow/compose.gleam").
-export([compose_sequential/3, compose_conditional/4, compose_parallel/4, validation_middleware/1]).
-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(" Composition API for combining flows.\n").
-file("src/telega/flow/compose.gleam", 188).
-spec create_composed_handler(
telega@flow@types:flow(gleam@dynamic:dynamic_(), AXVA, AXVB, AXVC),
list(telega@flow@types:flow(gleam@dynamic:dynamic_(), AXVA, AXVB, AXVC)),
integer()
) -> fun((telega@bot:context(AXVA, AXVB, AXVC), telega@flow@types:flow_instance()) -> {ok,
{telega@bot:context(AXVA, AXVB, AXVC),
telega@flow@types:flow_action(telega@flow@types:composed_step()),
telega@flow@types:flow_instance()}} |
{error, AXVB}).
create_composed_handler(Flow, All_flows, Index) ->
fun(Ctx, Instance) ->
{User_id, Chat_id} = telega@flow@engine:extract_ids_from_context(Ctx),
_pipe = telega@flow@engine:start_or_resume(
Flow,
Ctx,
User_id,
Chat_id,
erlang:element(3, erlang:element(6, Instance))
),
gleam@result:map(
_pipe,
fun(New_ctx) -> case gleam@list:drop(All_flows, Index + 1) of
[_ | _] ->
{New_ctx,
{next, {composed_flow_step, Index + 1}},
Instance};
[] ->
{New_ctx,
{complete,
erlang:element(3, erlang:element(6, Instance))},
Instance}
end end
)
end.
-file("src/telega/flow/compose.gleam", 159).
-spec string_to_composed_step(binary()) -> {ok,
telega@flow@types:composed_step()} |
{error, nil}.
string_to_composed_step(S) ->
case S of
<<"select_flow"/utf8>> ->
{ok, composed_select_flow};
<<"start_parallel"/utf8>> ->
{ok, composed_start_parallel};
<<"merge_results"/utf8>> ->
{ok, composed_merge_results};
_ ->
case gleam_stdlib:string_starts_with(S, <<"flow_"/utf8>>) of
true ->
_pipe = gleam@string:drop_start(S, 5),
_pipe@1 = gleam_stdlib:parse_int(_pipe),
_pipe@2 = gleam@result:map(
_pipe@1,
fun(Field@0) -> {composed_flow_step, Field@0} end
),
gleam@result:replace_error(_pipe@2, nil);
false ->
case gleam_stdlib:string_starts_with(
S,
<<"parallel_"/utf8>>
) of
true ->
_pipe@3 = gleam@string:drop_start(S, 9),
_pipe@4 = gleam_stdlib:parse_int(_pipe@3),
_pipe@5 = gleam@result:map(
_pipe@4,
fun(Field@0) -> {composed_parallel_flow, Field@0} end
),
gleam@result:replace_error(_pipe@5, nil);
false ->
{error, nil}
end
end
end.
-file("src/telega/flow/compose.gleam", 149).
-spec composed_step_to_string(telega@flow@types:composed_step()) -> binary().
composed_step_to_string(Step) ->
case Step of
{composed_flow_step, N} ->
<<"flow_"/utf8, (erlang:integer_to_binary(N))/binary>>;
composed_select_flow ->
<<"select_flow"/utf8>>;
composed_start_parallel ->
<<"start_parallel"/utf8>>;
{composed_parallel_flow, N@1} ->
<<"parallel_"/utf8, (erlang:integer_to_binary(N@1))/binary>>;
composed_merge_results ->
<<"merge_results"/utf8>>
end.
-file("src/telega/flow/compose.gleam", 23).
?DOC(" Compose flows sequentially\n").
-spec compose_sequential(
binary(),
list(telega@flow@types:flow(gleam@dynamic:dynamic_(), AXSV, AXSW, AXSX)),
telega@flow@types:flow_storage(AXSW)
) -> telega@flow@types:flow(telega@flow@types:composed_step(), AXSV, AXSW, AXSX).
compose_sequential(Name, Flows, Storage) ->
Flow_builder = telega@flow@builder:new(
Name,
Storage,
fun composed_step_to_string/1,
fun string_to_composed_step/1
),
Flow_builder@1 = gleam@list:index_fold(
Flows,
Flow_builder,
fun(B, Flow, Index) ->
Step = {composed_flow_step, Index},
telega@flow@builder:add_step(
B,
Step,
create_composed_handler(Flow, Flows, Index)
)
end
),
telega@flow@builder:build(Flow_builder@1, {composed_flow_step, 0}).
-file("src/telega/flow/compose.gleam", 41).
?DOC(" Compose flows with conditional selection\n").
-spec compose_conditional(
binary(),
fun((telega@flow@types:flow_instance()) -> binary()),
gleam@dict:dict(binary(), telega@flow@types:flow(gleam@dynamic:dynamic_(), AXTI, AXTJ, AXTK)),
telega@flow@types:flow_storage(AXTJ)
) -> telega@flow@types:flow(telega@flow@types:composed_step(), AXTI, AXTJ, AXTK).
compose_conditional(Name, Condition, Flows, Storage) ->
Flow_builder = telega@flow@builder:new(
Name,
Storage,
fun composed_step_to_string/1,
fun string_to_composed_step/1
),
Flow_builder@1 = telega@flow@builder:add_step(
Flow_builder,
composed_select_flow,
fun(Ctx, Instance) ->
Flow_name = Condition(Instance),
case gleam_stdlib:map_get(Flows, Flow_name) of
{ok, Flow} ->
{User_id, Chat_id} = telega@flow@engine:extract_ids_from_context(
Ctx
),
_pipe = telega@flow@engine:start_or_resume(
Flow,
Ctx,
User_id,
Chat_id,
erlang:element(3, erlang:element(6, Instance))
),
gleam@result:map(
_pipe,
fun(New_ctx) ->
{New_ctx,
{complete,
erlang:element(
3,
erlang:element(6, Instance)
)},
Instance}
end
);
{error, _} ->
{ok, {Ctx, cancel, Instance}}
end
end
),
telega@flow@builder:build(Flow_builder@1, composed_select_flow).
-file("src/telega/flow/compose.gleam", 75).
?DOC(" Compose flows for parallel execution\n").
-spec compose_parallel(
binary(),
list(telega@flow@types:flow(gleam@dynamic:dynamic_(), AXTW, AXTX, AXTY)),
fun((list(gleam@dict:dict(binary(), binary()))) -> gleam@dict:dict(binary(), binary())),
telega@flow@types:flow_storage(AXTX)
) -> telega@flow@types:flow(telega@flow@types:composed_step(), AXTW, AXTX, AXTY).
compose_parallel(Name, Flows, Merge_results, Storage) ->
Flow_builder = telega@flow@builder:new(
Name,
Storage,
fun composed_step_to_string/1,
fun string_to_composed_step/1
),
Parallel_steps = gleam@list:index_map(
Flows,
fun(_, Index) -> {composed_parallel_flow, Index} end
),
Flow_builder@1 = telega@flow@builder:add_step(
Flow_builder,
composed_start_parallel,
fun(Ctx, Instance) ->
{ok,
{Ctx,
{start_parallel, Parallel_steps, composed_merge_results},
Instance}}
end
),
Flow_builder@2 = gleam@list:index_fold(
Flows,
Flow_builder@1,
fun(B, Flow, Index@1) ->
Step = {composed_parallel_flow, Index@1},
telega@flow@builder:add_step(
B,
Step,
fun(Ctx@1, Instance@1) ->
{User_id, Chat_id} = telega@flow@engine:extract_ids_from_context(
Ctx@1
),
_pipe = telega@flow@engine:start_or_resume(
Flow,
Ctx@1,
User_id,
Chat_id,
erlang:element(3, erlang:element(6, Instance@1))
),
gleam@result:map(
_pipe,
fun(New_ctx) ->
{New_ctx,
{complete_parallel_step,
Step,
erlang:element(
3,
erlang:element(6, Instance@1)
)},
Instance@1}
end
)
end
)
end
),
Flow_builder@3 = telega@flow@builder:add_step(
Flow_builder@2,
composed_merge_results,
fun(Ctx@2, Instance@2) ->
case erlang:element(6, erlang:element(6, Instance@2)) of
{some, State} ->
Results = maps:values(erlang:element(4, State)),
Merged = Merge_results(Results),
Updated_instance = {flow_instance,
erlang:element(2, Instance@2),
erlang:element(3, Instance@2),
erlang:element(4, Instance@2),
erlang:element(5, Instance@2),
begin
_record = erlang:element(6, Instance@2),
{flow_state,
erlang:element(2, _record),
Merged,
erlang:element(4, _record),
erlang:element(5, _record),
erlang:element(6, _record)}
end,
erlang:element(7, Instance@2),
erlang:element(8, Instance@2),
erlang:element(9, Instance@2),
erlang:element(10, Instance@2),
erlang:element(11, Instance@2)},
{ok, {Ctx@2, {complete, Merged}, Updated_instance}};
none ->
{ok,
{Ctx@2,
{complete,
erlang:element(3, erlang:element(6, Instance@2))},
Instance@2}}
end
end
),
Flow_builder@4 = telega@flow@builder:add_parallel_steps(
Flow_builder@3,
composed_start_parallel,
Parallel_steps,
composed_merge_results
),
telega@flow@builder:build(Flow_builder@4, composed_start_parallel).
-file("src/telega/flow/compose.gleam", 133).
?DOC(" Validation middleware\n").
-spec validation_middleware(
fun((telega@flow@types:flow_instance()) -> {ok, nil} | {error, binary()})
) -> fun((telega@bot:context(AXUR, AXUS, AXUT), telega@flow@types:flow_instance(), fun(() -> {ok,
{telega@bot:context(AXUR, AXUS, AXUT),
telega@flow@types:flow_action(AXUQ),
telega@flow@types:flow_instance()}} |
{error, AXUS})) -> {ok,
{telega@bot:context(AXUR, AXUS, AXUT),
telega@flow@types:flow_action(AXUQ),
telega@flow@types:flow_instance()}} |
{error, AXUS}).
validation_middleware(Validator) ->
fun(Ctx, Instance, Next) -> case Validator(Instance) of
{ok, _} ->
Next();
{error, Msg} ->
case telega@reply:with_text(
Ctx,
<<"Validation failed: "/utf8, Msg/binary>>
) of
{ok, _} ->
{ok, {Ctx, back, Instance}};
{error, _} ->
{ok, {Ctx, cancel, Instance}}
end
end end.