Источник Задачи и gen_tcp
В этой главе мы изучим, как использовать модуль Erlang's :gen_tcp для обработки запросов. Это предоставляет отличную возможность изучить модуль Elixir's Task. В будущих главах мы расширим наш сервер, чтобы он мог обрабатывать команды.
Эхо-сервер
Мы начнем с реализации эхо-сервера. Он будет отправлять ответ с текстом, полученным в запросе. Мы будем постепенно улучшать сервер, пока он не будет контролироваться и готов обрабатывать несколько подключений.
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, так как нет смысла перезапускать не удавшееся подключение, но это плохой выбор для акцептора. Если акцептор аварийно завершается, мы хотим снова запустить акцептор.
Давайте исправим это. Мы знаем, что для дочернего элемента типа {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.16.3/task-and-gen-tcp.html