Packages
x3m_system
0.4.1
0.9.1
0.9.0
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
retired
0.7.20
0.7.19
0.7.18
0.7.17
0.7.16
0.7.15
0.7.14
0.7.13
0.7.12
0.7.11
0.7.10
0.7.9
0.7.8
retired
0.7.7
0.7.6
retired
0.7.5
0.7.4
retired
0.7.3
retired
0.7.2
0.7.1
0.7.0
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
retired
0.5.6
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.9
0.4.8
0.4.7
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.1.1
0.1.0
Building blocks for distributed and/or CQRS/ES systems
Current section
Files
Jump to
Current section
Files
lib/dispatcher.ex
defmodule X3m.System.Dispatcher do
alias X3m.System.{Message, Response, Instrumenter, ServiceRegistry}
def dispatch(%Message{halted?: true} = message), do: message
def dispatch(%Message{} = message, opts \\ []) do
timeout = Keyword.get(opts, :timeout, 5_000)
mono_start = System.monotonic_time()
Instrumenter.execute(
:discovering_service,
%{start: DateTime.utc_now(), mono_start: mono_start},
%{message: message, caller_node: Node.self()}
)
case discover_service(message) do
{:unavailable, message} ->
Instrumenter.execute(
:service_not_found,
%{time: DateTime.utc_now(), duration: Instrumenter.duration(mono_start, :microsecond)},
%{message: message, caller_node: Node.self()}
)
_unavailable(message)
{node, mod} ->
Instrumenter.execute(
:service_found,
%{time: DateTime.utc_now(), duration: Instrumenter.duration(mono_start, :microsecond)},
%{message: message, caller_node: Node.self(), service_node: node}
)
_dispatch(node, mod, message, timeout)
end
end
@spec discover_service(Message.t()) :: {:unavailable, Message.t()} | {node | :local, atom}
def discover_service(%Message{service_name: service} = message) do
case ServiceRegistry.find_nodes_with_service(service) do
:not_found -> {:unavailable, message}
{:local, {mod, _fun}} -> {:local, mod}
{:remote, nodes} -> Enum.random(nodes)
end
end
defp _dispatch(:local, mod, %Message{} = message, timeout) do
:ok = apply(mod, message.service_name, [message])
_wait_for_response(message, timeout)
end
defp _dispatch(node, mod, %Message{} = message, timeout) do
:ok = :rpc.call(node, mod, message.service_name, [message])
_wait_for_response(message, timeout)
end
defp _wait_for_response(message, timeout) do
receive do
%Message{} = message -> message
after
timeout ->
response = Response.service_timeout(message.service_name, message.id, timeout)
Message.return(message, response)
end
end
defp _unavailable(%Message{service_name: service_name} = message) do
response = Response.service_unavailable(service_name)
Message.return(message, response)
end
end