Spec-Zone.ru › Elixir 1.17

Исходный код Задача и gen_tcp

В этой главе мы изучим, как использовать модуль Erlang's :gen_tcp для обработки запросов. Это предоставляет отличную возможность изучить модуль Elixir's Task. В будущих главах мы расширим наш сервер, чтобы он мог обрабатывать команды.

Эхо-сервер

Мы начнем работу с TCP-сервером, реализовав эхо-сервер. Он будет отправлять ответ с текстом, полученным в запросе. Мы постепенно будем улучшать сервер, пока он не будет контролироваться и готов обрабатывать несколько подключений.

TCP-сервер, в общих чертах, выполняет следующие шаги:

  1. Прослушивает порт, пока порт не станет доступным, и получает доступ к сокету
  2. Ожидает подключение клиента на этом порту и принимает его
  3. Считывает запрос клиента и отправляет обратно ответ

Давайте реализуем эти шаги. Перейдите к приложению apps/kv_server, откройте lib/kv_server.ex, и добавьте следующие функции:

defmodule KVServer do
  require Logger

  def accept(port) do
    # The options below mean:
    #
    # 1. `:binary` - receives data as binaries (instead of lists)
    # 2. `packet: :line` - receives data line by line
    # 3. `active: false` - blocks on `:gen_tcp.recv/2` until data is available
    # 4. `reuseaddr: true` - allows us to reuse the address if the listener crashes
    #
    {:ok, socket} =
      :gen_tcp.listen(port, [:binary, packet: :line, active: false, reuseaddr: true])
    Logger.info("Accepting connections on port #{port}")
    loop_acceptor(socket)
  end

  defp loop_acceptor(socket) do
    {:ok, client} = :gen_tcp.accept(socket)
    serve(client)
    loop_acceptor(socket)
  end

  defp serve(socket) do
    socket
    |> read_line()
    |> write_line(socket)

    serve(socket)
  end

  defp read_line(socket) do
    {:ok, data} = :gen_tcp.recv(socket, 0)
    data
  end

  defp write_line(line, socket) do
    :gen_tcp.send(socket, line)
  end
end

Мы начнем работу с сервером, вызвав KVServer.accept(4040), где 4040 - это порт. Первый шаг в accept/1 - прослушивать порт, пока сокет не станет доступным, а затем вызвать loop_acceptor/1. loop_acceptor/1 - это цикл, принимающий подключения клиентов. Для каждого принятого подключения мы вызываем serve/1.

serve/1 - это другой цикл, который считывает строку из сокета и записывает эти строки обратно в сокет. Обратите внимание, что функция serve/1 использует оператор конвейера |>/2 для выражения этого потока операций. Оператор конвейера оценивает левую часть и передает результат в качестве первого аргумента функции в правой части.

socket |> read_line() |> write_line(socket)

эквивалентно:

write_line(read_line(socket), socket)

Реализация read_line/1 получает данные из сокета с помощью :gen_tcp.recv/2, а write_line/2 записывает в сокет с помощью :gen_tcp.send/2.

Обратите внимание, что serve/1 является бесконечным циклом, вызываемым последовательно внутри loop_acceptor/1, поэтому рекурсивный вызов loop_acceptor/1 никогда не достигается и может быть предотвращен. Однако, как мы увидим, нам потребуется выполнить serve/1 в отдельном процессе, поэтому нам скоро понадобится этот рекурсивный вызов.

Это в основном все, что нужно для реализации нашего эхо-сервера. Давайте попробуем!

Запустите сеанс IEx внутри приложения kv_server с помощью iex -S mix. Внутри IEx выполните:

iex> KVServer.accept(4040)

Сервер теперь запущен, и вы даже заметите, что консоль заблокирована. Давайте воспользуемся клиентом telnet для доступа к нашему серверу. Клиенты доступны на большинстве операционных систем, и их командные строки обычно похожи:

$ telnet 127.0.0.1 4040
Trying 127.0.0.1...
Connected to localhost.
Escape character is '^]'.
hello
hello
is it me
is it me
you are looking for?
you are looking for?

Введите "hello", нажмите Enter, и вы получите "hello" в ответ. Отлично!

Мой конкретный telnet-клиент может быть закрыт, набрав ctrl + ], набрав quit, и нажав <Enter>, но вашему клиенту могут потребоваться другие действия.

После выхода из telnet-клиента в сеансе IEx, вероятно, появится ошибка:

** (MatchError) no match of right hand side value: {:error, :closed}
    (kv_server) lib/kv_server.ex:45: KVServer.read_line/1
    (kv_server) lib/kv_server.ex:37: KVServer.serve/1
    (kv_server) lib/kv_server.ex:30: KVServer.loop_acceptor/1

Это потому, что мы ожидали данных от :gen_tcp.recv/2, но клиент закрыл соединение. В будущих версиях нашего сервера мы должны лучше обрабатывать такие случаи.

Пока, есть более важная ошибка, которую нужно исправить: что произойдет, если наш TCP-акцептор аварийно завершит работу? Поскольку нет контроля, сервер завершит работу, и мы не сможем обрабатывать больше запросов, так как он не будет перезапущен. Вот почему мы должны поместить наш сервер в дерево управления.

Задачи

Мы узнали об агентах, универсальных серверах и контроллерах. Все они предназначены для работы с несколькими сообщениями или управлением состоянием. Но что использовать, когда нам нужно просто выполнить какую-то задачу и всё?

Модуль Task предоставляет именно эту функциональность. Например, он имеет функцию Task.start_link/1, которая получает анонимную функцию и выполняет её в новом процессе, который будет частью дерева управления.

Давайте попробуем. Откройте lib/kv_server/application.ex, и измените контроллер в функции start/2 на следующее:

  def start(_type, _args) do
    children = [
      {Task, fn -> KVServer.accept(4040) end}
    ]

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

Как обычно, мы передали кортеж из двух элементов в качестве спецификации дочернего элемента, который, в свою очередь, вызовет Task.start_link/1.

С этим изменением мы говорим, что хотим запустить KVServer.accept(4040) как задачу. Мы жестко задаем порт на данный момент, но это можно изменить несколькими способами, например, прочитав порт из системной среды при запуске приложения:

port = String.to_integer(System.get_env("PORT") || "4040")
# ...
{Task, fn -> KVServer.accept(port) end}

Вставьте эти изменения в свой код, и теперь вы можете запустить своё приложение с помощью следующей команды PORT=4321 mix run --no-halt, обратите внимание, как мы передаём порт в качестве переменной, но по умолчанию используется 4040, если не указано иное.

Теперь, когда сервер является частью дерева управления, он должен запускаться автоматически при запуске приложения. Запустите сервер, теперь передавая порт, и еще раз используйте клиента telnet для проверки, что всё работает:

$ telnet 127.0.0.1 4321
Trying 127.0.0.1...
Connected to localhost.
Escape character is '^]'.
say you
say you
say me
say me

Да, всё работает! Однако, масштабируется ли это?

Попробуйте подключить два telnet-клиента одновременно. Когда вы это сделаете, вы заметите, что второй клиент не откликается:

$ telnet 127.0.0.1 4321
Trying 127.0.0.1...
Connected to localhost.
Escape character is '^]'.
hello
hello?
HELLOOOOOO?

По всей видимости, это вообще не работает. Это потому, что мы обрабатываем запросы в том же процессе, что и принимаем подключения. Когда подключен один клиент, мы не можем принять другого.

Контроллер задач

Чтобы наш сервер обрабатывал одновременные подключения, нам нужно, чтобы один процесс работал в качестве акцептора, который запускает другие процессы для обработки запросов. Одним из решений было бы изменить:

defp loop_acceptor(socket) do
  {:ok, client} = :gen_tcp.accept(socket)
  serve(client)
  loop_acceptor(socket)
end

чтобы также использовать Task.start_link/1:

defp loop_acceptor(socket) do
  {:ok, client} = :gen_tcp.accept(socket)
  Task.start_link(fn -> serve(client) end)
  loop_acceptor(socket)
end

Мы запускаем связанную задачу непосредственно из процесса акцептора. Но мы уже совершали эту ошибку раньше. Вы помните?

Это похоже на ошибку, которую мы допустили, когда вызвали KV.Bucket.start_link/1 напрямую из реестра. Это означало, что отказ любого контейнера приводил к отказу всего реестра.

В приведенном выше коде будет такой же недостаток: если мы связали задачу serve(client) с акцептором, отказ при обработке запроса приведет к отказу акцептора и, следовательно, всех других подключений.

Мы исправили проблему для реестра, используя простой один-к-одному контроллер. Мы будем использовать ту же тактику здесь, за исключением того, что этот шаблон настолько распространён с задачами, что Task уже поставляется с решением: простой один-к-одному контроллер, который запускает временные задачи как часть нашего дерева управления.

Давайте изменим start/2 еще раз, чтобы добавить контроллер в наше дерево:

  def start(_type, _args) do
    port = String.to_integer(System.get_env("PORT") || "4040")

    children = [
      {Task.Supervisor, name: KVServer.TaskSupervisor},
      {Task, fn -> KVServer.accept(port) end}
    ]

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

Теперь мы запустим процесс Task.Supervisor с именем KVServer.TaskSupervisor. Помните, поскольку задача акцептора зависит от этого контроллера, контроллер должен быть запущен первым.

Теперь нам нужно изменить loop_acceptor/1 на использование Task.Supervisor для обработки каждого запроса:

defp loop_acceptor(socket) do
  {:ok, client} = :gen_tcp.accept(socket)
  {:ok, pid} = Task.Supervisor.start_child(KVServer.TaskSupervisor, fn -> serve(client) end)
  :ok = :gen_tcp.controlling_process(client, pid)
  loop_acceptor(socket)
end

Вы можете заметить, что мы добавили строку :ok = :gen_tcp.controlling_process(client, pid). Это делает дочерний процесс "управляющим процессом" сокета client. Если бы мы этого не сделали, акцептор привел бы к падению всех клиентов при аварийном завершении работы, потому что сокеты были бы привязаны к процессу, который их принял (что является стандартным поведением).

Запустите новый сервер с PORT=4040 mix run --no-halt и теперь вы можете открыть несколько одновременных telnet-клиентов. Вы также заметите, что завершение работы клиента не приводит к падению акцептора. Отлично!

Вот полная реализация эхо-сервера:

defmodule KVServer do
  require Logger

  @doc """
  Starts accepting connections on the given `port`.
  """
  def accept(port) do
    {:ok, socket} = :gen_tcp.listen(port,
                      [:binary, packet: :line, active: false, reuseaddr: true])
    Logger.info "Accepting connections on port #{port}"
    loop_acceptor(socket)
  end

  defp loop_acceptor(socket) do
    {:ok, client} = :gen_tcp.accept(socket)
    {:ok, pid} = Task.Supervisor.start_child(KVServer.TaskSupervisor, fn -> serve(client) end)
    :ok = :gen_tcp.controlling_process(client, pid)
    loop_acceptor(socket)
  end

  defp serve(socket) do
    socket
    |> read_line()
    |> write_line(socket)

    serve(socket)
  end

  defp read_line(socket) do
    {:ok, data} = :gen_tcp.recv(socket, 0)
    data
  end

  defp write_line(line, socket) do
    :gen_tcp.send(socket, line)
  end
end

Поскольку мы изменили спецификацию контроллера, мы должны спросить: наша стратегия управления по-прежнему верна?

В этом случае ответ да: если акцептор аварийно завершает работу, нет необходимости аварийно завершать существующие соединения. С другой стороны, если контроллер задач аварийно завершает работу, нет необходимости аварийно завершать работу акцептора тоже.

Однако, остается одна проблема, которая связана со стратегиями перезапуска. У задач по умолчанию значение :restart установлено в :temporary, что означает, что они не перезапускаются. Это отличное значение по умолчанию для соединений, инициированных через Task.Supervisor, так как перезапуск завершенного соединения не имеет смысла, но это плохой выбор для акцептора. Если акцептор аварийно завершает работу, мы хотим перезапустить его.

Давайте исправим это. Мы знаем, что для дочернего элемента формы {Task, fun} Elixir вызовет Task.child_spec(fun) для получения базовой спецификации дочернего элемента. Следовательно, можно предположить, что для изменения спецификации {Task, fun} на имеющую :restart значение :permanent, нам нужно изменить модуль Task. Однако это невозможно, так как модуль Task определен как часть стандартной библиотеки Elixir (и даже если это было бы возможно, маловероятно, что это была бы хорошая идея). К счастью, это можно сделать, используя Supervisor.child_spec/2, который позволяет настроить спецификацию дочернего элемента с новыми значениями. Давайте перепишем start/2 в KVServer.Application еще раз:

  def start(_type, _args) do
    port = String.to_integer(System.get_env("PORT") || "4040")

    children = [
      {Task.Supervisor, name: KVServer.TaskSupervisor},
      Supervisor.child_spec({Task, fn -> KVServer.accept(port) end}, restart: :permanent)
    ]

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

Теперь у нас есть постоянно работающий акцептор, который запускает временные процессы задач под постоянно работающим контроллером задач.

В следующей главе мы начнем парсить запросы клиентов и отправлять ответы, завершая работу нашего сервера.

← Предыдущая страница Зависимости и проекты-оболочки
Следующая страница → Тесты документации, шаблоны и с

Скачать версию ePub

Создано с использованием ExDoc (v0.34.1) для язык программирования Elixir

© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.17.2/task-and-gen-tcp.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API