Исходный код Задачи и gen_tcp
В этой главе мы изучим, как использовать модуль Erlang :gen_tcp для обработки запросов. Это отличная возможность изучить модуль Elixir Task. В будущих главах мы расширим наш сервер, чтобы он мог обрабатывать команды.
Эхо-сервер
Мы начнем с реализации эхо-сервера TCP. Он будет отправлять ответ, содержащий текст, полученный в запросе. Мы будем постепенно улучшать сервер, пока он не станет контролируемым и готовым обрабатывать множество подключений.
TCP-сервер, в общих чертах, выполняет следующие шаги:
- Прослушивает порт, пока порт не станет доступным и не получит сокет
- Ожидает подключение клиента на этом порте и принимает его
- Читает запрос клиента и отправляет ответ обратно
Давайте реализуем эти шаги. Перейдите к приложению 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, так как нет смысла перезапускать failed подключение, но это плохой выбор для акцептора. Если акцептор выйдет из строя, мы хотим перезапустить его.
Давайте это исправим. Мы знаем, что для дочернего элемента типа {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
Теперь у нас есть постоянно работающий акцептор, который запускает временные задачи под управлением постоянно работающего надзирателя задач.
В следующей главе мы начнём разбирать запросы клиентов и отправлять ответы, завершая создание нашего сервера.
© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.18.1/task-and-gen-tcp.html