~/adam.log

Concurrent Data Processing in Elixir - Chapter 2

Published 2026-10-03

2. Long-Running Processes Using GenServer

So what we have learned so far is a quick and easy way to build a concurrent set of instruction and use functions. These small sets of processes deal with trash clean-up and don’t require much on the server side.

What happens when we need a more robust and longer lasting structure? Well that is where GenServer comes in. We will start to work with GenServers in this chapter.


Starting with a Basic GenServer

We will continue to work with the sender project for this next set of code. First create a new file in the lib folder called send_server.ex we can then define the module like this.

defmodule SendServer do
  use GenServer
end

We can start the program with the normal start we used before

iex -S mix

With the output

Erlang/OTP 27 [erts-15.1.2] [source] [64-bit] [smp:24:24] [ds:24:24:10] [async-threads:1] [jit:ns]

Compiling 1 file (.ex)
    warning: function init/1 required by behaviour GenServer is not implemented (in module SendServer).

    We will inject a default implementation for now:

        def init(init_arg) do
          {:ok, init_arg}
        end

    You can copy the implementation above or define your own that converts the arguments given to GenServer.start_link/3 to the server state.

    │
  1 │ defmodule SendServer do
    │ ~~~~~~~~~~~~~~~~~~~~~~~
    │
    └─ lib/send_server.ex:1: SendServer (module)

Generated sender app
Interactive Elixir (1.18.0) - press Ctrl+C to exit (type h() ENTER for help)
iex(1)> 

You did it you have the first GenServer


GenServer Callbacks In Depth

So we will now need to talk about callbacks that will be used to do all the work within the GenServer. there are defaults that we will be overriding to get the functionality that we want. In order to override something we need to know:

  1. What arguments the callback takes
  2. What the return values that are supported

We will cover the following within this chapter:

  • handle_call/3
  • handle_cast/2
  • handle_continue/2
  • handle_info/2
  • init/1
  • terminate/2

Learning just these will get you very far into understanding how the GenServer functions.

Initializing the Process

Let’s start with init/1 this is callback that is started when you first start the process. We will now start to make it so the SendServer can see faux emails.

  @impl true
  def init(args) do
    IO.puts("Received arguments: #{inspect(args)}")
    max_retries = Keyword.get(args, :max_retries, 5)
    state = %{emails: [], max_retries: max_retries}
    {:ok, state}
  end

We set up some standard setting within the init and we can now use it as we need. Head back to your iex session and test it out.

iex(3)> {:ok, pid} = GenServer.start(SendServer, [max_retries: 1])
Received arguments: [max_retries: 1]
{:ok, #PID<0.187.0>}
iex(4)> GenServer.stop(pid)
:ok

There are some standard returns for the init some common ones are the following:

}
:ignore

We covered the but the is great for doing post-initialization work.

Lets now use the handle_continue/2 to deal with some post init work to get more information. As its not good to do much besides set the initial state of the process within the init/1.

  @impl true
  def handle_continue(:fetch_from_database, state) do
    # called after init/1
  end

Breaking Down Work in Multiple Steps

Okay so some of the return values for handle_continue/2 are:

}

We can use one of those now with this implementation. It will assume that we have a DB that has a users

  @impl true
  def handle_continue(:fetch_from_database, state) do
    users = []
    # get `users` from the database
    {:noreply, Map.put(state, :users, users)}
  end

Sending Process Messages

One of the nice things about this GenServer is that you can send and receive messages from the process while its running that is where GenServer.call/3 and GenServer.cast/2 will come in. You can call and expect a reply and you can cast and send info without a need for a reply. Let’s add in some more code.

  @impl true
  def handle_call(:get_state, _from, state) do
    {:reply, state, state}
  end

We have said that we want to get a reply and set the state as the reply, notice that you also need to send the state in the return meaning that you can do work on the state within a call. Here are some of the normal returns:

}

Now we can recompile the iex session and test it out.

iex(7)> {:ok, pid} = GenServer.start(SendServer, [max_retries: 1])
Received arguments: [max_retries: 1]
{:ok, #PID<0.214.0>}
iex(8)> GenServer.call(pid, :get_state)
%{emails: [], max_retries: 1}

Okay that is great now we can implement sending emails with the handle_cast/2

  def handle_cast({:send, email}, state) do
    Sender.send_email(email)
    emails = [%{email: email, status: "sent", retries: 0}] ++ state.emails
    {:noreply, %{state | emails: emails}}
  end

This assumes that you remember that we implemented the Sender.send_email earlier in the book. We can now try that out.

iex(10)> {:ok, pid} = GenServer.start(SendServer, [max_retries: 1])
Received arguments: [max_retries: 1]
{:ok, #PID<0.232.0>}
iex(11)> GenServer.cast(pid, {:send, "hello@email.com"})
:ok
iex(12)> GenServer.call(pid, :get_state)
Email to hello@email.com sent
%{
  emails: [%{status: "sent", email: "hello@email.com", retries: 0}],
  max_retries: 1
}

Notifying the Process of Events

Now we can talk about an other callback that we can implement the handle_info/2 this is meant to deal with times that you get new information. We will use the Process.send/2 to send info to the Process, it will then trigger the handle_info/2 callback. Let’s work with it now. We can even make it have to work through its retries. Augment the send_email/1 logic to the following and we can augment the handle_cast/2 as well.

  def send_email("konnichiwa@world.com" = email),
    do: :error

  def send_email(email) do
    Process.sleep(3000)
    IO.puts("Email to #{email} sent")
    {:ok, "email_sent"}
  end


  @impl true
  def handle_cast({:send, email}, state) do
    status =
      case Sender.send_email(email) do
        {:ok, "email_sent"} -> "sent"
        :error -> "failed"
      end

    emails = [%{email: email, status: status, retries: 0}] ++ state.emails
    {:noreply, %{state | emails: emails}}
  end

Okay so we have a way to deal with errors and a way to retry. Let’s add a line to the init/1 so that we try to send out an email asap.

Process.send_after(self(), :retry, 5000)

Now let’s setup the handle_info/2

  @impl true
  def handle_info(:retry, state) do
    {failed, done} =
      Enum.split_with(state.emails, fn item ->
        item.status == "failed" && item.retries < state.max_retries
      end)

    retried =
      Enum.map(failed, fn item ->
        IO.puts("Retrying email #{item.email}...")

        new_status =
          case Sender.send_email(item.email) do
            {:ok, "email_sent"} -> "sent"
            :error -> "failed"
          end

        %{email: item.email, status: new_status, retries: item.retries + 1}
      end)

    Process.send_after(self(), :retry, 5000)
    {:noreply, %{state | emails: retried ++ done}}
  end

We now have a way of retrying the emails after we start. It will continue to do so over and over as in the beginning it will send a retry, then end of the handle_info/2 it will send out an other :retry. Let’s test it out in the iex session.

iex(15)> {:ok, pid} = GenServer.start(SendServer, max_retries: 2)
Received arguments: [max_retries: 2]
{:ok, #PID<0.257.0>}
iex(16)> GenServer.cast(pid, {:send, "hello@world.com"})
:ok
iex(17)> GenServer.cast(pid, {:send, "aloha@world.com"})
:ok
iex(18)> GenServer.cast(pid, {:send, "konnichiwa@world.com"})
:ok
Email to aloha@world.com sent
iex(20)> GenServer.call(pid, :get_state)
%{
  emails: [
    %{status: "failed", email: "konnichiwa@world.com", retries: 0},
    %{status: "sent", email: "aloha@world.com", retries: 0},
    %{status: "sent", email: "hello@world.com", retries: 0}
  ],
  max_retries: 2
}

Process Teardown

Okay there is an other callback that we should talk about and its the terminate it is what is run when you kill a process. This is a small bit of code but add it to the Module.

  @impl true
  def terminate(reason, _state) do
    IO.puts("Terminating with reason #{reason}")
  end

Let’s test it out.

iex(22)> {:ok, pid} = GenServer.start(SendServer, [])
Received arguments: []
{:ok, #PID<0.276.0>}
iex(23)> GenServer.stop(pid)
Terminating with reason normal
:ok

Building a Job-Processing System

Okay so we want to build a process server that will scale to make it run a process for every email that we want to send, with all the functionality that we had already the retries and so on. We will start a new project for this so let’s get started.

mix new jobber --sup

We can set up a new file job.ex in lib/jobber with the following:

defmodule Jobber.Job do
  use GenServer
  require Logger
end

Initializing the Job Process

Lets set-up the init/1 for the new job.ex

  defstruct [:work, :id, :max_retries, retries: 0, status: "new"]

  def init(args) do
    work = Keyword.fetch!(args, :work)
    id = Keyword.get(args, :id, random_job_id())
    max_retries = Keyword.get(args, :max_retries, 3)
    state = %Jobber.Job{id: id, work: work, max_retries: max_retries}
    {:ok, state, {:continue, :run}}
  end

  defp random_job_id() do
    :crypto.strong_rand_bytes(5) |> Base.url_encode64(padding: false)
  end

We want to have a unique :id for each job so we will use the :crypto to do so. Because we added the need for :crypto we need to add it into the mix.exs

  def application do
    [
      extra_applications: [:logger, :crypto],
      mod: {Jobber.Application, []}
    ]
  end

Performing Work

We set the return to be } so we will need to handle that continue right away in-order to finish the init/1

  @impl true
  def handle_continue(:run, state) do
    new_state = state.work.() |> handle_job_result(state)

    if new_state.status == "errored" do
      Process.send_after(self(), :retry, 5000)
      {:noreply, new_state}
    else
      Logger.info("Job exiting #{state.id}")
      {:stop, :normal, new_state}
    end
  end

Okay so now we need a way to deal with the different states of the work. So we will need to define some private functions with pattern matching. Based off these states:

  • Success, when the job completes and returns {:ok, data}
  • Initial error, when it fails the first time with :error
  • Retry error, when we attempt to rerun the job and also receive :error
  defp handle_job_result({:ok, _data}, state) do
    Logger.info("Job completed #{state.id}")
    %Jobber.Job{state | status: "done"}
  end

  defp handle_job_result(:error, %{status: "new"} = state) do
    Logger.warn("Job errored #{state.id}")
    %Jobber.Job{state | status: "errored"}
  end

  defp handle_job_result(:error, %{status: "errored"} = state) do
    Logger.warn("Job retry failed #{state.id}")
    new_state = %Jobber.Job{state | retries: state.retries + 1}

    if new_state.retries == state.max_retries do
      %Jobber.Job{new_state | status: "failed"}
    else
      new_state
    end
  end

Now we can handle the handle_info/2 for the :retry message that will be sent.

  @impl true
  def handle_info(:retry, state) do
    # Delegate work to the `handle_continue/2` callback.
    {:noreply, state, {:continue, :run}}
  end

You should be able to see that we will go back to the continue branch of the logic and that will start it over again. Let’s test it out.

iex(1)> GenServer.start(Jobber.Job, work: fn -> Process.sleep(5000); {:ok, []} end)
{:ok, #PID<0.188.0>}
iex(2)> 
nil
iex(3)> good_job = fn ->
...(3)>   Process.sleep(5000)
...(3)>   {:ok, []}
...(3)> end
#Function<43.39164016/0 in :erl_eval.expr/6>
iex(4)> 
nil
iex(5)> GenServer.start(Jobber.Job, work: good_job)
{:ok, #PID<0.189.0>}
iex(6)> 
nil
iex(7)> bad_job = fn ->
...(7)>   Process.sleep(5000)
...(7)>   :error
...(7)> end
#Function<43.39164016/0 in :erl_eval.expr/6>
iex(8)> 
nil
iex(9)> GenServer.start(Jobber.Job, work: bad_job)
{:ok, #PID<0.190.0>}
iex(10)> 
nil

19:35:56.155 [info] Job completed uOdxGQg

19:35:56.155 [warning] Job errored cQ0pp-s

19:35:56.158 [info] Job exiting uOdxGQg

19:35:56.155 [info] Job completed -xPCT84

19:35:56.158 [info] Job exiting -xPCT84

Introducing DynamicSupervisor

Okay so before we had everything did start again but it needed to be started with a GenServer.start, that starts the process without a link-process that will die with the child and the process will be lost. There is a better way and that is the DynamicSupervisor. We will need to head into the application.ex in-order to create the supervisor tree for this.

  @impl true
  def start(_type, _args) do
    children = [
      {DynamicSupervisor, strategy: :one_for_one, name: Jobber.JobRunner}
    ]

    # See https://hexdocs.pm/elixir/Supervisor.html
    # for other strategies and supported options
    opts = [strategy: :one_for_one, name: Jobber.Supervisor]
    Supervisor.start_link(children, opts)
  end
end

Let’s now create the .iex.exs so that we don’t have to add in the good and bad jobs.

good_job = fn ->
  Process.sleep(5000)
  {:ok, []}
end

bad_job = fn ->
  Process.sleep(5000)
  :error
end

Okay so we have the functions now let’s start to work with the jobber.ex this will be what we use to start jobs within the DynamicSupervisor.

#jobber/lib/jobber.ex
defmodule Jobber do
  alias Jobber.{JobRunner, Job}

  def start_job(args) do
    DynamicSupervisor.start_child(JobRunner, {Job, args})
  end
end

Let’s test it.

iex(1)> Jobber.start_job(work: good_job)
{:error,
 {:undef,
  [
    {Jobber.Job, :start_link,
     [[work: #Function<43.39164016/0 in :erl_eval.expr/6>]], []},
    {DynamicSupervisor, :start_child, 3,
     [file: ~c"lib/dynamic_supervisor.ex", line: 800]},
    {DynamicSupervisor, :handle_start_child, 2,
     [file: ~c"lib/dynamic_supervisor.ex", line: 786]},
    {:gen_server, :try_handle_call, 4, [file: ~c"gen_server.erl", line: 2381]},
    {:gen_server, :handle_msg, 6, [file: ~c"gen_server.erl", line: 2410]},
    {:proc_lib, :init_p_do_apply, 3, [file: ~c"proc_lib.erl", line: 329]}
  ]}}

Looks like we don’t have a start_link/2 let’s get that defined.

#jobber/lib/jobber/job.ex
def start_link(args) do
  GenServer.start_link(__MODULE__, args)
end

Now it works.

iex(4)> Jobber.start_job(work: good_job)
{:ok, #PID<0.182.0>}

Process Restart Values Revisited

There are ways to make it so it will not restart if the process exits normally. Otherwise it will restart forever. Replace the macro for the GenServer with

use GenServer, restart: :transient

Now lets test the 2 different jobs.

iex(1)> Jobber.start_job(work: good_job)
{:ok, #PID<0.164.0>}

19:51:13.954 [info] Job completed og7u9bA

19:51:13.956 [info] Job exiting og7u9bA

iex(2)> Jobber.start_job(work: bad_job)
{:ok, #PID<0.165.0>}

19:52:15.013 [warning] Job errored b16AJNU

19:52:25.017 [warning] Job retry failed b16AJNU

19:52:32.155 [warning] Job retry failed b16AJNU

19:52:42.167 [warning] Job retry failed b16AJNU

19:52:42.167 [info] Job exiting b16AJNU

That all worked as expected but what happens when we have a doomed job? Add the following to the .iex.exs

doomed_job = fn ->
  Process.sleep(5000)
  raise "Boom!"
end
iex(1)> Jobber.start_job(work: doomed_job)
{:ok, #PID<0.150.0>}

19:54:29.974 [error] GenServer #PID<0.150.0> terminating
** (RuntimeError) Boom!
    (elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
    (jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
    (stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
    (stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "FjeMdjg", max_retries: 3, retries: 0, status: "new"}

19:54:35.014 [error] GenServer #PID<0.151.0> terminating
** (RuntimeError) Boom!
    (elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
    (jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
    (stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
    (stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "60TqQzk", max_retries: 3, retries: 0, status: "new"}

Adjusting Restart Frequency

Right now we have the doomed job because the timeout is longer than the frequency can handle. With the max_restarts set to 3 and the max_seconds set to 5. We can change those. Go to the application.ex and add this extra setting.

  @impl true
  def start(_type, _args) do
    job_runner_config = [
      strategy: :one_for_one,
      max_seconds: 30,
      name: Jobber.JobRunner
    ]

    children = [
      {DynamicSupervisor, job_runner_config}
    ]

    opts = [strategy: :one_for_one, name: Jobber.Supervisor]
    Supervisor.start_link(children, opts)
  end

Lets test it out and keep track of the pids.

iex(1)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
iex(2)> Jobber.start_job(work: doomed_job)
{:ok, #PID<0.163.0>}
iex(3)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>

19:59:06.705 [error] GenServer #PID<0.163.0> terminating
** (RuntimeError) Boom!
    (elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
    (jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
    (stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
    (stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "6U5INHM", max_retries: 3, retries: 0, status: "new"}

19:59:11.731 [error] GenServer #PID<0.164.0> terminating
** (RuntimeError) Boom!
    (elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
    (jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
    (stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
    (stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "K_c4Xmg", max_retries: 3, retries: 0, status: "new"}
iex(4)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>

19:59:16.740 [error] GenServer #PID<0.165.0> terminating
** (RuntimeError) Boom!
    (elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
    (jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
    (stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
    (stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "jUeDQis", max_retries: 3, retries: 0, status: "new"}
iex(5)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>

19:59:21.747 [error] GenServer #PID<0.166.0> terminating
** (RuntimeError) Boom!
    (elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
    (jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
    (stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
    (stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "-FI_W5E", max_retries: 3, retries: 0, status: "new"}
iex(6)> Process.whereis(Jobber.JobRunner)
#PID<0.167.0>

Looking at the output we see that when a process fails the attached process is killed and we are left with a new Jobber. We can make sure that the same process is kept with a Supervisor.


Implementing a Supervisor

Both of the 2 Supervisors that we have reached for are good for simple trees but we want to have more control over the Supervisors and the children. So now we can reach for The Supervisor module. We can start by creating a new file job_supervisor.ex.

defmodule Jobber.JobSupervisor do
  use Supervisor, restart: :temporary

  def start_link(args) do
    Supervisor.start_link(__MODULE__, args)
  end

  def init(args) do
    children = [
      {Jobber.Job, args}
    ]

    options = [
      strategy: :one_for_one,
      max_seconds: 30
    ]

    Supervisor.init(children, options)
  end
end

Now that we have the new Supervisor we can change the Jobber to use the new Supervisor

defmodule Jobber do
  alias Jobber.{JobRunner, Job, JobSupervisor}

  def start_job(args) do
    DynamicSupervisor.start_child(JobRunner, {JobSupervisor, args})
  end
end

Let’s test it out.

iex(1)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
iex(2)> Jobber.start_job(work: doomed_job)
{:ok, #PID<0.163.0>}

20:10:12.559 [error] GenServer #PID<0.163.0> terminating
** (RuntimeError) Boom!
    (elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
    (jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
    (stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
    (stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "mvdyZfs", max_retries: 3, retries: 0, status: "new"}
iex(3)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>

20:10:17.601 [error] GenServer #PID<0.164.0> terminating
** (RuntimeError) Boom!
    (elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
    (jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
    (stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
    (stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "jJD80-Y", max_retries: 3, retries: 0, status: "new"}
iex(4)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>

See the same PID?

Okay so understand what we have here there are 3 types of restart structure that you can use:

  • :one_for_one
  • :one_for_all
  • :rest_for_one

Depending on which one you use, everything will be restarted, nothing will be restarted or pieces will be restarted.

Naming Processes Using the Registry

Okay so one of the most powerful options is the ability to name a process to call it by name, we have been using the pid identifier so far but we can change that.

children = [
{DynamicSupervisor, strategy: :one_for_one, name: Jobber.JobRunner},
]

Typically, the :name is set when we link a process using start_link/1. There are three types of accepted values when naming a process:

  • An atom, like :job_runner. This includes module names, since they’re atoms under the hood, for example, Jobber.JobRunner.
  • A {:global, term} tuple, like {:global, :job_runner}, which registers the process globally. Useful for distributed applications
  • A {:via, module, term} tuple, where module is an Elixir module that would take care of the registration process, using the value term

Starting a Registry Process

This is actually a very simple process head back to the application.ex and make the following changes.

  @impl true
  def start(_type, _args) do
    job_runner_config = [
      strategy: :one_for_one,
      max_seconds: 30,
      name: Jobber.JobRunner
    ]

    children = [
      # This next line
      {Registry, keys: :unique, name: Jobber.JobRegistry},
      {DynamicSupervisor, job_runner_config}
    ]

    opts = [strategy: :one_for_one, name: Jobber.Supervisor]
    Supervisor.start_link(children, opts)
  end

Registering New Processes

Okay so we can now use the Registry to name a process, lets head to the job.ex and make a few changes starting by adding in a via helper.

  defp via(key, value) do
    {:via, Registry, {Jobber.JobRegistry, key, value}}
  end

  def start_link(args) do
    args =
      if Keyword.has_key?(args, :id) do
        args
      else
        Keyword.put(args, :id, random_job_id())
      end

    id = Keyword.get(args, :id)
    type = Keyword.get(args, :type)
    GenServer.start_link(__MODULE__, args, name: via(id, type))
  end  

    @impl true
  def init(args) do
    work = Keyword.fetch!(args, :work)
    # This line next
    id = Keyword.get(args, :id)
    max_retries = Keyword.get(args, :max_retries, 3)
    state = %Jobber.Job{id: id, work: work, max_retries: max_retries}
    {:ok, state, {:continue, :run}}
  end

Querying the Registry

Registry allows you to use the select/2 it might look cryptic but once you get to understand it you should be fine.

  def running_imports() do
    match_all = {:"$1", :"$2", :"$3"}
    guards = [{:==, :"$3", "import"}]
    map_result = [%{id: :"$1", pid: :"$2", type: :"$3"}]
    Registry.select(Jobber.JobRegistry, [{match_all, guards, map_result}])
  end

The match specification we have is a bit lengthy, so we have broken it down to three variables:

  • match_all is a wildcard that matches all entries in the registry.
  • guards is a filter that filters results by the third element in the tuple, which has to be equal to “import”.
  • map_result is transforming the result by creating a list of maps, assigning each element of the tuple to a key, which makes the result a bit more readable.

Now go to the .iex.exs and increase the sleep for the good job.

good_job = fn ->
  Process.sleep(60_000)
  {:ok, []}
end

So now we have a way of querying the Registry. Let’s test it out.

iex(1)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.178.0>}
iex(2)> Jobber.start_job(work: good_job, type: "send_email")
{:ok, #PID<0.179.0>}
iex(3)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.180.0>}
iex(4)> Jobber.running_imports()
[
  %{id: "xjzkd1k", pid: #PID<0.178.0>, type: "import"},
  %{id: "MaTRVTw", pid: #PID<0.180.0>, type: "import"}
]
iex(5)> 

Limiting Concurrency of Important Jobs

Now we can really use the running_imports/0 to good use. Add this check to Jobber.start_job/1

  alias Jobber.{JobRunner, Job, JobSupervisor}

  # ...

  def start_job(args) do
    if Enum.count(running_imports()) >= 5 do
      {:error, :import_quota_reached}
    else
      DynamicSupervisor.start_child(JobRunner, {JobSupervisor, args})
    end
  end

Now test it out by trying to add more than 5 jobs.

iex(1)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.166.0>}
iex(2)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.168.0>}
iex(3)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.170.0>}
iex(4)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.172.0>}
iex(5)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.174.0>}
iex(6)> Jobber.start_job(work: good_job, type: "import")
{:error, :import_quota_reached}

Inspecting Supervisors at Runtime

You can query a running supervision tree with two functions. Both are available in Supervisor and DynamicSupervisor.

count_children/1

Takes a supervisor’s PID or name and returns a map of child counts:

  • :supervisors counts child supervisors, both active and inactive.
  • :workers counts child workers, both active and inactive.
  • :specs counts all children, both active and inactive.
  • :active counts the children that are running right now.
iex(1)> DynamicSupervisor.count_children(Jobber.JobRunner)
%{active: 0, specs: 0, supervisors: 0, workers: 0}

which_children/1

Returns a list of {id, pid, type, [module]} tuples:

  • id is :undefined for dynamically started children.
  • pid is the child’s PID, or :restarting while it is being restarted.
  • type is :worker or :supervisor.
  • [module] is the implementing module, inside a list.

The Idle JobSupervisor

After a job finishes, JobRunner still shows active: 1, supervisors: 1.

iex(2)> DynamicSupervisor.which_children(Jobber.JobRunner)
[{:undefined, #PID<0.152.0>, :supervisor, [Jobber.JobSupervisor]}]

The Job process has exited, but its parent JobSupervisor keeps running with nothing to do. To confirm this, check the JobSupervisor directly:

iex(3)> {_, pid, _, _} = List.first(DynamicSupervisor.which_children(Jobber.JobRunner))
iex(4)> Supervisor.count_children(pid)
%{active: 0, specs: 1, supervisors: 0, workers: 1}

Idle processes are cheap and usually harmless. If you want to clean one up, stop it:

iex(5)> Supervisor.stop(pid)
:ok

After that, DynamicSupervisor.which_children(Jobber.JobRunner) returns [].

Use count_children/1 to get the numbers and which_children/1 to see what each child is. Together they let you find leftover processes in the tree and stop them.

Wrapping Up

We now have the supervisor trees set and you can work with them.

I would take some time now and get used to doing more and more with the tree. You can work on persistency and even adding more into queues.