-module(temporal_sdk_api_common). % elp:ignore W0012 W0040 -moduledoc false. -export([ retry_policy/1, get_by_casted_key/4, mutation_reset_reason/1, marker_encode_value/2, marker_decode_value/2, run_request/6, run_request/7, transform_workflow_execution/1, put_run_id/3, format_response/4, t2e/3 ]). -include("proto.hrl"). -spec retry_policy(RetryPolicy :: temporal_sdk:retry_policy()) -> ?TEMPORAL_SPEC:'temporal.api.common.v1.RetryPolicy'(). retry_policy(RetryPolicy) -> RP1 = case RetryPolicy of #{initial_interval := II} -> RetryPolicy#{initial_interval := temporal_sdk_utils_time:msec_to_protobuf(II)}; #{} -> RetryPolicy end, case RP1 of #{maximum_interval := MI} -> RP1#{maximum_interval := temporal_sdk_utils_time:msec_to_protobuf(MI)}; #{} -> RP1 end. -spec get_by_casted_key( Key :: term(), Map :: map(), Default :: term(), ApiCtx :: temporal_sdk_api:context() ) -> {CastedKey :: unicode:chardata(), Value :: term()}. get_by_casted_key(Key, Map, Default, #{client_opts := #{grpc_opts := #{codec := {Codec, _, _}}}}) -> CastedKey = Codec:cast(Key), case Map of #{CastedKey := V} -> {CastedKey, V}; #{} -> {CastedKey, Default} end. -spec mutation_reset_reason(Name :: unicode:chardata() | atom()) -> {string(), binary()}. mutation_reset_reason(Name) -> Str = temporal_sdk_utils_path:string_path([?MUTATION_RESET_PREFIX, Name], ":"), {Str, temporal_sdk_utils_unicode:characters_to_binary1(Str)}. -spec marker_encode_value( Codec :: temporal_sdk_workflow:marker_value_codec(), Value :: temporal_sdk:term_to_payloads() | term() ) -> {ok, EncodedValue :: temporal_sdk:term_to_payloads(), Decoder :: none | list | dlist | term | fun((temporal_sdk:term_from_payloads()) -> term()) | {Module :: module(), Function :: atom()}} | {error, Reason :: map()}. marker_encode_value(list, Value) -> {ok, [Value], list}; marker_encode_value(dlist, Value) -> {ok, [[Value]], list}; marker_encode_value(term, Value) -> {ok, [erlang:term_to_binary(Value)], term}; marker_encode_value(none, Value) when is_list(Value) -> {ok, Value, none}; marker_encode_value({none, none}, Value) when is_list(Value) -> {ok, Value, none}; marker_encode_value({none, DF}, Value) when is_list(Value), is_function(DF, 1) -> {ok, Value, DF}; marker_encode_value({none, {DM, DF}}, Value) when is_list(Value), is_atom(DM), is_atom(DF) -> {ok, Value, {DM, DF}}; marker_encode_value({EF, none}, Value) when is_function(EF, 1) -> case EF(Value) of EV when is_list(EV) -> {ok, EV, none}; EV -> do_encode_type_err(EF, Value, EV) end; marker_encode_value({{EM, EF}, none}, Value) when is_atom(EM), is_atom(EF) -> case EM:EF(Value) of EV when is_list(EV) -> {ok, EV, none}; EV -> do_encode_type_err({EM, EF}, Value, EV) end; marker_encode_value({EF, none}, Value) when is_function(EF, 1) -> case EF(Value) of EV when is_list(EV) -> {ok, EV, none}; EV -> do_encode_type_err(EF, Value, EV) end; marker_encode_value({EF, DF}, Value) when is_function(EF, 1), is_function(DF, 1) -> case EF(Value) of EV when is_list(EV) -> {ok, EV, none}; EV -> do_encode_type_err(EF, Value, EV) end; marker_encode_value({EF, {DM, DF}}, Value) when is_function(EF, 1), is_atom(DM), is_atom(DF) -> case EF(Value) of EV when is_list(EV) -> {ok, EV, {DM, DF}}; EV -> do_encode_type_err(EF, Value, EV) end; marker_encode_value(Codec, Value) -> {error, #{ reason => "Invalid marker value type or invalid marker configuration.", marker_value => Value, marker_value_codec => Codec }}. do_encode_type_err(Encoder, Value, EncodedValue) -> {error, #{ reason => "Invalid encoded marker value type. Expected a list() type.", marker_value => Value, invalid_encoded_marker_value => EncodedValue, encoder => Encoder }}. -spec marker_decode_value( Decoder :: none | list | dlist | term | fun((temporal_sdk:term_from_payloads()) -> term()) | {Module :: module(), Function :: atom()}, Value :: temporal_sdk:term_from_payloads() ) -> DecodedValue :: term(). marker_decode_value(list, [Value]) -> Value; marker_decode_value(dlist, [[Value]]) -> Value; marker_decode_value(term, [Value]) -> erlang:binary_to_term(Value); marker_decode_value(none, Value) -> Value; marker_decode_value(DF, Value) when is_function(DF, 1) -> DF(Value); marker_decode_value({DM, DF}, Value) when is_atom(DM), is_atom(DF) -> DM:DF(Value). -spec run_request( Cluster :: temporal_sdk_cluster:cluster_name(), Opts :: proplists:proplist(), DefaultOpts :: temporal_sdk_utils_opts:defaults(), ServiceName :: temporal_sdk_api:temporal_service(), RequestMessageName :: temporal_sdk_client:msg_name(), ResponseMessageName :: temporal_sdk_client:msg_name() ) -> temporal_sdk:response(). run_request(Cluster, Opts, DefaultOpts, ServiceName, RequestMessageName, ResponseMessageName) -> run_request( Cluster, Opts, DefaultOpts, ServiceName, RequestMessageName, ResponseMessageName, [] ). -spec run_request( Cluster :: temporal_sdk_cluster:cluster_name(), Opts :: proplists:proplist(), DefaultOpts :: temporal_sdk_utils_opts:defaults(), ServiceName :: temporal_sdk_api:temporal_service(), RequestMessageName :: temporal_sdk_client:msg_name(), ResponseMessageName :: temporal_sdk_client:msg_name(), Customizations :: proplists:proplist() ) -> temporal_sdk:response(). run_request( Cluster, Opts, DefaultOpts, ServiceName, RequestMessageName, ResponseMessageName, Customizations ) -> maybe {ok, ApiCtx} ?= temporal_sdk_api_context:build(Cluster), {ok, FullOpts} ?= temporal_sdk_utils_opts:build(DefaultOpts, Opts, ApiCtx), ok ?= maybe_evict(Cluster, FullOpts, proplists:to_map(Customizations)), {RawRequest, O1} = maps:take(raw_request, FullOpts), {ResponseType, O2} = maps:take(response_type, O1), {GrpcOpts, O3} = maps:take(grpc_opts, O2), Req1 = maps:merge(O3, RawRequest), Req = do_custom(Customizations, ApiCtx, RequestMessageName, Req1), Response = temporal_sdk_api:request(ServiceName, Cluster, Req, ResponseType, GrpcOpts), format_response(ResponseMessageName, ResponseType, Response, ApiCtx) end. maybe_evict( Cluster, #{namespace := Namespace}, #{evict := EvictReason, workflow_execution := WorkflowExecution} ) -> WE = case WorkflowExecution of {_, W} -> W; W -> W end, temporal_sdk:evict_workflow(Cluster, WE, [{namespace, Namespace}, {reason, EvictReason}]); maybe_evict(_Cluster, _Opts, _Customizations) -> ok. do_custom([{new, {Key, Val}} | TC], ApiCtx, MsgName, Req) -> do_custom(TC, ApiCtx, MsgName, Req#{Key => Val}); do_custom([{nested, {PKey, CKeys}} | TC], ApiCtx, MsgName, Req) -> R = maps:without(CKeys, Req), do_custom(TC, ApiCtx, MsgName, R#{PKey => maps:with(CKeys, Req)}); do_custom([{workflow_execution, {Key, WE}} | TC], ApiCtx, MsgName, Req) -> do_custom(TC, ApiCtx, MsgName, Req#{Key => transform_workflow_execution(WE)}); do_custom([{workflow_execution, WE} | TC], ApiCtx, MsgName, Req) -> do_custom([{workflow_execution, {workflow_execution, WE}} | TC], ApiCtx, MsgName, Req); do_custom([identity | TC], ApiCtx, MsgName, Req) -> do_custom(TC, ApiCtx, MsgName, temporal_sdk_api:put_identity(ApiCtx, MsgName, Req)); do_custom([id | TC], ApiCtx, MsgName, Req) -> do_custom([{id, id} | TC], ApiCtx, MsgName, Req); do_custom([{id, Key} | TC], ApiCtx, MsgName, Req) -> do_custom(TC, ApiCtx, MsgName, temporal_sdk_api:put_id(ApiCtx, MsgName, Key, Req)); do_custom([{evict, _} | TC], ApiCtx, MsgName, Req) -> do_custom(TC, ApiCtx, MsgName, Req); do_custom([], _ApiCtx, _MsgName, Req) -> Req. -spec transform_workflow_execution( WorkflowExecutionOrId :: temporal_sdk:workflow_execution_or_id() ) -> temporal_sdk:workflow_execution(). transform_workflow_execution(#{workflow_id := _} = WE) -> WE; transform_workflow_execution(WId) when is_binary(WId); is_list(WId) -> #{workflow_id => WId}. -spec put_run_id( Cluster :: temporal_sdk_cluster:cluster_name(), WorkflowExecutionOrId :: temporal_sdk:workflow_execution_or_id(), Opts :: map() ) -> {ok, temporal_sdk:workflow_execution()} | {error, Reason :: term()}. put_run_id(_Cluster, #{run_id := _} = WE, _Opts) -> {ok, WE}; put_run_id(Cluster, WorkflowExecutionOrId, #{namespace := Namespace}) -> case temporal_sdk:describe_workflow(Cluster, WorkflowExecutionOrId, [{namespace, Namespace}]) of {ok, #{workflow_execution_info := #{execution := #{run_id := RunId}}}} -> WE = temporal_sdk_api_common:transform_workflow_execution(WorkflowExecutionOrId), {ok, WE#{run_id => RunId}}; {ok, Response} -> {error, #{ reason => "Invalid describe_workflow response - missing run_id.", invalid_response => Response }}; Err -> % eqwalizer:ignore Err end. -spec format_response( MessageName :: temporal_sdk_client:msg_name(), ResponseType :: temporal_sdk:response_type(), Response :: {ok, temporal_sdk_client:msg()} | {error, term()}, ApiCtx :: temporal_sdk_api:context() ) -> temporal_sdk:response(). format_response(MessageName, call_formatted, {ok, Response}, ApiCtx) -> {ok, t2e(MessageName, Response, ApiCtx)}; format_response( _MessageName, call_formatted, {error, #{grpc_response_headers := #{~"grpc-message" := Err}}}, _ApiCtx ) -> {error, Err}; format_response( _MessageName, call_formatted, {error, #{grpc_response_headers := #{"grpc-message" := Err}}}, _ApiCtx ) -> {error, Err}; format_response(_MessageName, _Opts, Response, _ApiCtx) -> Response. -spec t2e( MessageName :: temporal_sdk_client:msg_name(), Message :: temporal_sdk_client:msg(), ApiCtx :: temporal_sdk_api:context() ) -> ConvertedMessage :: map(). t2e(MessageName, Message, ApiCtx) when is_map(Message) -> maps:from_list(te(MessageName, maps:to_list(Message), ApiCtx, [], [])). te(MN, [{K, #{payloads := _} = V} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [ {K, temporal_sdk_api:map_from_payloads(ApiCtx, MN, KAcc ++ [K], V)} | Acc ]); te(MN, [{K, #{data := _, metadata := _} = V} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [ {K, temporal_sdk_api:map_from_payload(ApiCtx, MN, KAcc ++ [K], V)} | Acc ]); te(MN, [{header = K, #{fields := V}} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [ {K, temporal_sdk_api:map_from_mapstring_payload(ApiCtx, MN, KAcc ++ [K, fields], V)} | Acc ]); te(MN, [{memo = K, #{fields := V}} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [ {K, temporal_sdk_api:map_from_mapstring_payload(ApiCtx, MN, KAcc ++ [K, fields], V)} | Acc ]); te(MN, [{search_attributes = K, #{indexed_fields := V}} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [ {K, temporal_sdk_api:map_from_mapstring_payload(ApiCtx, MN, KAcc ++ [K, fields], V)} | Acc ]); te(MN, [{workflow_type = K, #{name := N}} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [{K, N} | Acc]); te(MN, [{task_queue = K, #{name := N}} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [{K, N} | Acc]); te(MN, [{type = K, #{name := N}} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [{K, N} | Acc]); te(MN, [{K, #{seconds := _, nanos := _} = T} | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [{K, temporal_sdk_utils_time:protobuf_to_msec(T)} | Acc]); te(MN, [{K, [#{} | _] = Li} | TM], ApiCtx, KAcc, Acc) -> V = lists:map( fun(M) when is_map(M) -> maps:from_list(te(MN, maps:to_list(M), ApiCtx, KAcc ++ [K], [])) end, Li ), te(MN, TM, ApiCtx, KAcc, [{K, V} | Acc]); te(MN, [{K, #{} = Ma} | TM], ApiCtx, KAcc, Acc) -> V = maps:from_list(te(MN, maps:to_list(Ma), ApiCtx, KAcc ++ [K], [])), te(MN, TM, ApiCtx, KAcc, [{K, V} | Acc]); te(MN, [{K, {KT, #{} = Ma}} | TM], ApiCtx, KAcc, Acc) -> V = maps:from_list(te(MN, maps:to_list(Ma), ApiCtx, KAcc ++ [K], [])), te(MN, TM, ApiCtx, KAcc, [{K, {KT, V}} | Acc]); te(MN, [KV | TM], ApiCtx, KAcc, Acc) -> te(MN, TM, ApiCtx, KAcc, [KV | Acc]); te(_MN, [], _ApiCtx, _KAcc, Acc) -> Acc.