defmodule Kvasir.Agent.Instance do use GenServer require Logger alias Kvasir.Offset @startup_timeout 60_000 @keep_alive 15 * 60_000 @keep_hibernated 5 * @keep_alive defp log_start(start_result, agent, id) defp log_start(err = {:error, reason}, agent, id) when is_atom(reason) do Logger.error("Agent[#{id}]: Failed to start. (reason: #{inspect(reason)}", agent: agent, id: id, reason: reason ) err end defp log_start(start_result, _agent, _id), do: start_result def start_agent(agent, partition, id, opts) do opts = Keyword.put(opts, :timeout, @startup_timeout) config = :config |> agent.__agent__() |> Map.put(:id, id) |> Map.put(:partition, partition) |> Map.update!(:cache, &elem(&1, 0)) __MODULE__ |> GenServer.start_link(config, opts) |> log_start(agent, id) end def pre_start_agent(agent, partition, id, offset, state, cache, opts \\ []) do opts = Keyword.put(opts, :timeout, @startup_timeout) config = %{model: m} = agent.__agent__(:config) with {:ok, s} <- m.__decode__(state) do c = config |> Map.put(:id, id) |> Map.put(:partition, partition) |> Map.put(:agent_state, s) |> Map.put(:offset, offset) |> Map.put(:cache, cache) |> Map.put(:callbacks, %{}) __MODULE__ |> GenServer.start_link({:preload, c}, opts) |> log_start(agent, id) end end @impl GenServer def init({:preload, state = %{agent: agent, partition: partition}}) do start = :erlang.monotonic_time() Logger.debug(fn -> "Agent<#{state.id}>: Preloaded-Init (#{inspect(self())})" end) initialize(state) |> report_init(start, agent, partition) end def init( state = %{ source: source, agent: agent, partition: partition, id: id, cache: cache_mod, topic: topic, model: model } ) do start = :erlang.monotonic_time() Logger.debug(fn -> "Agent<#{state.id}>: Init (#{inspect(self())})" end) with {:ok, c} <- cache_mod.cache(agent, partition, id) do cache = {cache_mod, c} state = state |> Map.put(:cache, cache) |> Map.put(:callbacks, %{}) case load_state(source, cache, topic, model, id) do {:ok, {offset, agent_state}} -> state |> Map.put(:agent_state, agent_state) |> Map.put(:offset, offset) |> initialize() :rebuild -> {:ok, state, {:continue, :rebuild}} end end |> report_init(start, agent, partition) end @doc false @impl GenServer def handle_continue(continue, state) def handle_continue( :rebuild, state = %{ source: source, id: id, cache: {cache_m, cache_i}, topic: topic, model: model } ) do with {:ok, {offset, agent_state}} <- build_state(source, topic, model, id, nil, model.base(id)) do cache_m.save(cache_i, model.__encode__(agent_state), offset) state |> Map.put(:agent_state, agent_state) |> Map.put(:offset, offset) |> initialize(:noreply) end end @spec initialize(map, success :: :ok | :noreply) :: {:ok, map} defp initialize(state = %{agent: agent}, success \\ :ok) do Logger.debug(fn -> "Agent<#{state.id}>: Initialized" end) keep_alive = if @keep_alive != :infinity, do: Process.send_after(self(), {:hibernate, :keep_alive}, @keep_alive) publisher = Module.concat([agent, "Client"]) {success, state |> Map.put_new(:publisher, &publisher.publish/1) |> Map.put(:hibernate, nil) |> Map.put(:keep_alive, keep_alive)} end defp load_state(_source, {cache_m, cache_i}, _topic, model, id) do case cache_m.load(cache_i) do :no_previous_state -> Logger.debug(fn -> "Agent<#{id}>: State Loaded: No State, Setup New" end) {:ok, {Kvasir.Offset.create(), model.base(id)}} {:ok, offset, state} -> Logger.debug(fn -> "Agent<#{id}>: State Loaded: #{inspect(offset)}" end) # build_state(source, topic, model, id, offset, state) with {:ok, s} <- model.__decode__(state), do: {:ok, {offset, s}} {:error, :corrupted_state} -> Logger.warn(fn -> "Agent<#{id}>: State Load Failed: Corrupted State" end) :rebuild {:error, reason} -> Logger.error(fn -> "Agent<#{id}>: State Load Failed: #{reason}" end) :rebuild end end defp build_state(source, topic, model, id, offset, original_state) do base = {offset || Kvasir.Offset.create(), original_state, false} with {:ok, {a, b, _}} <- topic |> source.stream(from: offset, to: :last, key: id) |> Stream.filter(&is_map/1) |> EnumX.reduce_while(base, &state_reducer(model, &1, &2)), do: {:ok, {a, b}} end defp state_reducer(_model, err = {:error, _}, _), do: err defp state_reducer(model, event, {offset, state, true}) do with {:ok, updated_state} <- model.apply(state, event) do {:ok, {Offset.set(offset, event.__meta__.partition, event.__meta__.offset), updated_state, true}} end end defp state_reducer(model, event, {offset, state, false}) do if Offset.empty?(offset) or Offset.get(offset, event.__meta__.partition) < event.__meta__.offset do state_reducer(model, event, {offset, state, true}) else {:ok, {offset, state, false}} end end defp state_reducer(_model, _event, err), do: err @impl GenServer def handle_info( {:command, from, command}, state = %{ agent: agent, agent_state: agent_state, cache: {cache_m, cache_i}, model: model, partition: p } ) do start = :erlang.monotonic_time() pause_keep_alive(state) Logger.debug(fn -> "Agent<#{state.id}>: Command: #{inspect(command)}" end) ref = command.__meta__.id response = case model.execute(agent_state, command) do {:ok, events} when is_list(events) -> :ok = cache_m.track_command(cache_i) report_execute(true, start, agent, p) commit_events(state, events, ref) {:ok, event} -> :ok = cache_m.track_command(cache_i) report_execute(true, start, agent, p) commit_events(state, [event], ref) :ok -> report_execute(true, start, agent, p) :ok error -> report_execute(false, start, agent, p) error end r = case response do {:ok, events} -> send(from, {:command, ref, {:ok, Offset.create(p, List.last(events).__meta__.offset)}}) with {:ok, state} <- apply_events(events, state), do: {:noreply, reset_keep_alive(state)} _ -> send(from, {:command, ref, response}) {:noreply, reset_keep_alive(state)} end report_command(response, start, agent, p) r end def handle_info({:event, event}, state) do with {:ok, state} <- apply_events([event], state), do: {:noreply, state} end def handle_info({:offset_callback, from, ref, offset}, state) do Logger.debug(fn -> "Agent<#{state.id}>: Adding callback (#{inspect(offset)})" end) {:noreply, add_offset_callback(state, {from, ref}, offset)} end def handle_info({:hibernate, reason}, state) do Logger.debug(fn -> "Agent<#{state.id}>: Hibernate (#{inspect(reason)})" end) hibernate = if @keep_hibernated != :infinity, do: Process.send_after(self(), {:shutdown, :hibernated}, @keep_hibernated) {:noreply, %{state | hibernate: hibernate}, :hibernate} end def handle_info({:shutdown, reason}, state) do Logger.debug(fn -> "Agent<#{state.id}>: Shutdown (#{inspect(reason)})" end) {:stop, :shutdown, state} end def handle_info({:inet_reply, _, _}, state), do: {:noreply, state} @impl GenServer def handle_call(:inspect, _from, state) do {:reply, state.agent_state, state} end def handle_call(:rebuild, _from, state) do case rebuild(state) do {:ok, s} -> {:reply, :ok, s} err -> {:reply, err, state} end end @impl GenServer def handle_cast(request, state) def handle_cast({:rebuild, callback}, state) do case rebuild(state) do {:ok, s} -> callback.(:ok) {:noreply, s} err -> callback.(err) {:noreply, state} end end defp rebuild(state) do %{source: source, cache: {cache_m, cache_i}, topic: topic, model: model, id: id} = state Logger.debug(fn -> "Agent<#{id}>: Forced State Rebuild" end) with {:ok, {offset, agent_state}} <- build_state(source, topic, model, id, nil, model.base(id)) do cache_m.save(cache_i, model.__encode__(agent_state), offset) updated_state = state |> Map.put(:agent_state, agent_state) |> Map.put(:offset, offset) {:ok, updated_state} else {:ok, {_, _, false}} -> {:ok, state} err -> err end end defp apply_events(events, state = %{agent: agent, partition: partition}) do start = :erlang.monotonic_time() with {:ok, new_state = %{offset: o}, updated} <- do_apply_events(events, state) do if updated do %{cache: {cache_m, cache_i}, model: model, agent_state: agent_state} = new_state :ok = cache_m.save(cache_i, model.__encode__(agent_state), o) end notify_offset_callbacks(new_state, o) {:ok, new_state} end |> report_apply(start, agent, partition) end defp do_apply_events(events, state, updated \\ false) defp do_apply_events([], state, updated), do: {:ok, state, updated} defp do_apply_events([event | events], state = %{offset: o}, updated) do Logger.debug(fn -> "Agent<#{state.id}>: Incoming Event (#{inspect(event)})" end) %{offset: offset, partition: partition} = event.__meta__ if Offset.empty?(o) or Offset.get(o, event.__meta__.partition) < offset do case state.model.apply(state.agent_state, event) do :ok -> new_state = %{state | offset: Offset.set(state.offset, partition, offset)} do_apply_events(events, new_state, true) {:ok, updated_state} -> updated_offset = Offset.set(state.offset, partition, offset) new_state = %{state | offset: updated_offset, agent_state: updated_state} do_apply_events(events, new_state, true) {:error, reason} -> if Kvasir.Event.on_error(event) == :halt do Logger.error(fn -> "Agent<#{state.id}>: Event error (#{inspect(reason)})" end) {:stop, :invalid_event, state} else Logger.warn(fn -> "Agent<#{state.id}>: Event error (#{inspect(reason)})" end) # Should this be true? We updated by skipping I guess do_apply_events(events, state, true) end end else do_apply_events(events, state, updated) end end defp pause_keep_alive(%{hibernate: h_ref, keep_alive: a_ref}) do h_ref && Process.cancel_timer(h_ref) a_ref && Process.cancel_timer(a_ref) :ok end defp reset_keep_alive(state = %{keep_alive: nil}), do: state defp reset_keep_alive(state) do pause_keep_alive(state) %{ state | hibernate: nil, keep_alive: Process.send_after(self(), {:hibernate, :keep_alive}, @keep_alive) } end defp commit_events(state = %{publisher: publish}, events, ref) do events |> Enum.map(&prepare_event(&1, ref, state)) |> EnumX.map(&publish.(&1)) end defp prepare_event(event, _ref, %{id: id, partition: partition}) do # meta = %{event.__meta__ | command: ref, partition: partition, key: id} meta = %{event.__meta__ | partition: partition, key: id} %{event | __meta__: meta} end defp add_offset_callback(state = %{callbacks: callbacks, offset: now}, l = {pid, ref}, offset) do if Offset.compare(now, offset) == :lt do callbacks = Map.update(callbacks, offset, [l], &[l | &1]) %{state | callbacks: callbacks} else send(pid, {:offset_reached, ref, offset}) state end end def notify_offset_callbacks(state = %{callbacks: callbacks}, offset) do grouped = Enum.group_by(callbacks, &(Offset.compare(elem(&1, 0), offset) != :lt)) Enum.each(Map.get(grouped, true, []), fn {off, listeners} -> Enum.each(listeners, fn {pid, ref} -> send(pid, {:offset_reached, ref, off}) end) end) %{state | callbacks: Enum.into(Map.get(grouped, false, []), %{})} end @compile {:inline, success?: 1, report_init: 4, report_execute: 4, report_command: 4, report_apply: 4, report: 5} @spec report_init(term, integer, module, integer) :: term defp report_init(result, start, agent, partition), do: report(result, start, agent, partition, "kvasir.agent.initialize.timer:") @spec report_execute(boolean, integer, module, integer) :: :ok defp report_execute(result, start, agent, partition) do report( if(result, do: {:ok, true}, else: false), start, agent, partition, "kvasir.agent.execute.timer:" ) :ok end @spec report_command(boolean, integer, module, integer) :: :ok defp report_command(result, start, agent, partition), do: report(result, start, agent, partition, "kvasir.agent.command.timer:") @spec report_apply(term, integer, module, integer) :: term defp report_apply(result, start, agent, partition), do: report(result, start, agent, partition, "kvasir.agent.apply.timer:") @spec report(term, integer, module, integer, String.t()) :: term defp report(result, start, agent, partition, label) do stop = :erlang.monotonic_time() ms = :erlang.convert_time_unit(stop - start, :native, :millisecond) Kvasir.Agent.Metrics.Sender.send([ label, to_string(ms), "|ms|#agent:", inspect(agent), ",partition:", to_string(partition), success?(result) ]) result end defp success?(result) defp success?(:ok), do: ",success:true" defp success?({:ok, _}), do: ",success:true" defp success?(_), do: ",success:false" end