TCPConnection
TCP-соединение. При подключении используется алгоритм Happy Eyeballs.
Следующий код создаёт клиента, который подключается к порту 8989 локального хоста, записывает «hello world» и ожидает ответа, который затем выводит.
use "net"
class MyTCPConnectionNotify is TCPConnectionNotify
let _out: OutStream
new create(out: OutStream) =>
_out = out
fun ref connected(conn: TCPConnection ref) =>
conn.write("hello world")
fun ref received(
conn: TCPConnection ref,
data: Array[U8] iso,
times: USize)
: Bool
=>
_out.print("GOT:" + String.from_array(consume data))
conn.close()
true
fun ref connect_failed(conn: TCPConnection ref) =>
None
actor Main
new create(env: Env) =>
try
TCPConnection(env.root as AmbientAuth,
recover MyTCPConnectionNotify(env.out) end, "", "8989")
end
Примечание: при записи в соединение данные будут молча проигнорированы, если соединение ещё не установлено.
Поддержка обратной пропускной способности
Запись
Протокол TCP имеет встроенную поддержку обратной пропускной способности. Это обычно проявляется в том, что буфер исходящей записи заполняется и все запрошенные данные не могут быть записаны в сокет. В TCPConnection, это скрыто от программиста. В этом случае TCPConnection будет буферизовать дополнительные данные до тех пор, пока они не смогут быть отправлены. Если это не контролировать, может возникнуть неконтролируемое очередирование. Для решения этой проблемы TCPConnectionNotify реализует два метода throttled и unthrottled, которые вызываются при применении и освобождении обратной пропускной способности.
При получении уведомления throttled у вашего приложения есть два варианта обработки. Один из них — сообщить среде выполнения Pony, что она больше не может делать прогресс, и что должна быть применена обратная пропускная способность среды выполнения к любым акторам, отправляющим сообщения этому актору. Например, вы можете создать своё приложение так:
// Here we have a TCPConnectionNotify that upon construction
// is given a BackpressureAuth token. This allows the notifier
// to inform the Pony runtime when to apply and release backpressure
// as the connection experiences it.
// Note the calls to
//
// Backpressure.apply(_auth)
// Backpressure.release(_auth)
//
// that apply and release backpressure as needed
use "backpressure"
use "collections"
use "net"
class SlowDown is TCPConnectionNotify
let _auth: BackpressureAuth
let _out: StdStream
new iso create(auth: BackpressureAuth, out: StdStream) =>
_auth = auth
_out = out
fun ref throttled(connection: TCPConnection ref) =>
_out.print("Experiencing backpressure!")
Backpressure.apply(_auth)
fun ref unthrottled(connection: TCPConnection ref) =>
_out.print("Releasing backpressure!")
Backpressure.release(_auth)
fun ref closed(connection: TCPConnection ref) =>
// if backpressure has been applied, make sure we release
// when shutting down
_out.print("Releasing backpressure if applied!")
Backpressure.release(_auth)
fun ref connect_failed(conn: TCPConnection ref) =>
None
actor Main
new create(env: Env) =>
try
let auth = env.root as AmbientAuth
let socket = TCPConnection(auth, recover SlowDown(auth, env.out) end,
"", "7669")
end
Или, если хотите, вы можете обработать обратную пропускную способность, снизив нагрузку, то есть, отказавшись от дополнительных данных, вместо выполнения отправки. Это может выглядеть так:
use "net"
class ThrowItAway is TCPConnectionNotify
var _throttled: Bool = false
fun ref sent(conn: TCPConnection ref, data: ByteSeq): ByteSeq =>
if not _throttled then
data
else
""
end
fun ref sentv(conn: TCPConnection ref, data: ByteSeqIter): ByteSeqIter =>
if not _throttled then
data
else
recover Array[String] end
end
fun ref throttled(connection: TCPConnection ref) =>
_throttled = true
fun ref unthrottled(connection: TCPConnection ref) =>
_throttled = false
fun ref connect_failed(conn: TCPConnection ref) =>
None
actor Main
new create(env: Env) =>
try
TCPConnection(env.root as AmbientAuth,
recover ThrowItAway end, "", "7669")
end
В общем случае, если у вас нет очень специфического случая использования, мы настоятельно рекомендуем не реализовывать схему снижения нагрузки, где вы отбрасываете данные.
Чтение
Если ваше приложение не может справляться с данными, отправляемыми в него по TCPConnection, вы можете использовать встроенную поддержку обратной пропускной способности при чтении, чтобы приостановить чтение сокета, что, в свою очередь, начнёт оказывать обратную пропускную способность на соответствующего писателя на другом конце сокета.
Поведение mute позволяет любым другим акторам в вашем приложении запросить прекращение дополнительных чтений до тех пор, пока не будет вызван unmute. Обратите внимание, что это прекращение не гарантируется, так как это результат асинхронного вызова поведения, и поэтому ему придётся подождать, пока все сообщения в почтовом ящике TCPConnection не будут обработаны.
На платформах, отличных от Windows, ваше TCPConnection не заметит, если другой конец соединения закрыт, пока вы его не разбудите. Системы Unix, такие как FreeBSD, Linux и OSX, узнают о закрытом соединении при чтении. На этих платформах вы обязательно должны вызвать unmute для заглушенного соединения, чтобы оно закрылось. Без вызова unmute актёр TCPConnection никогда не завершится.
Поддержка прокси
Используя обратный вызов proxy_via в TCPConnectionNotify можно реализовать прокси-серверы. Функция принимает целевой хост и службу в качестве параметров и возвращает 2-кортеж из хоста и службы прокси.
Прокси-сервер TCPConnectionNotify должен украшать другую реализацию TCPConnectionNotify и передавать соответствующие данные.
Пример реализации прокси
actor Main
new create(env: Env) =>
MyClient.create(
"example.com", // we actually want to connect to this host
"80",
ExampleProxy.create("proxy.example.com", "80")) // we connect via this proxy
actor MyClient
new create(host: String, service: String, proxy: Proxy = NoProxy) =>
let conn: TCPConnection = TCPConnection.create(
env.root as AmbientAuth,
proxy.apply(MyConnectionNotify.create()),
host,
service)
class ExampleProxy is Proxy
let _proxy_host: String
let _proxy_service: String
new create(proxy_host: String, proxy_service: String) =>
_proxy_host = proxy_host
_proxy_service = proxy_service
fun apply(wrap: TCPConnectionNotify iso): TCPConnectionNotify iso^ =>
ExampleProxyNotify.create(consume wrap, _proxy_service, _proxy_service)
class iso ExampleProxyNotify is TCPConnectionNotify
// Fictional proxy implementation that has no error
// conditions, and always forwards the connection.
let _proxy_host: String
let _proxy_service: String
var _destination_host: (None | String) = None
var _destination_service: (None | String) = None
let _wrapped: TCPConnectionNotify iso
new iso create(wrap: TCPConnectionNotify iso, proxy_host: String, proxy_service: String) =>
_wrapped = wrap
_proxy_host = proxy_host
_proxy_service = proxy_service
fun ref proxy_via(host: String, service: String): (String, String) =>
// Stash the original host & service; return the host & service
// for the proxy; indicating that the initial TCP connection should
// be made to the proxy
_destination_host = host
_destination_service = service
(_proxy_host, _proxy_service)
fun ref connected(conn: TCPConnection ref) =>
// conn is the connection to the *proxy* server. We need to ask the
// proxy server to forward this connection to our intended final
// destination.
conn.write((_destination_host + "\n").array())
conn.write((_destination_service + "\n").array())
wrapped.connected(conn)
fun ref received(conn, data, times) => _wrapped.received(conn, data, times)
fun ref connect_failed(conn: TCPConnection ref) => None
actor tag TCPConnection
Конструкторы
create
Подключение через IPv4 или IPv6. Если from — непустая строка, подключение будет установлено с указанного интерфейса.
new tag create(
auth: (AmbientAuth val | NetAuth val | TCPAuth val |
TCPConnectAuth val),
notify: TCPConnectionNotify iso,
host: String val,
service: String val,
from: String val = "",
read_buffer_size: USize val = 16384,
yield_after_reading: USize val = 16384,
yield_after_writing: USize val = 16384)
: TCPConnection tag^
Параметры
- auth: (AmbientAuth val | NetAuth val | TCPAuth val | TCPConnectAuth val)
- notify: TCPConnectionNotify iso
- host: String val
- service: String val
- from: String val = ""
- read_buffer_size: USize val = 16384
- yield_after_reading: USize val = 16384
- yield_after_writing: USize val = 16384
Возвращаемое значение
- TCPConnection tag^
ip4
Подключение через IPv4.
new tag ip4(
auth: (AmbientAuth val | NetAuth val | TCPAuth val |
TCPConnectAuth val),
notify: TCPConnectionNotify iso,
host: String val,
service: String val,
from: String val = "",
read_buffer_size: USize val = 16384,
yield_after_reading: USize val = 16384,
yield_after_writing: USize val = 16384)
: TCPConnection tag^
Параметры
- auth: (AmbientAuth val | NetAuth val | TCPAuth val | TCPConnectAuth val)
- notify: TCPConnectionNotify iso
- host: String val
- service: String val
- from: String val = ""
- read_buffer_size: USize val = 16384
- yield_after_reading: USize val = 16384
- yield_after_writing: USize val = 16384
Возвращаемое значение
- TCPConnection tag^
ip6
Подключение через IPv6.
new tag ip6(
auth: (AmbientAuth val | NetAuth val | TCPAuth val |
TCPConnectAuth val),
notify: TCPConnectionNotify iso,
host: String val,
service: String val,
from: String val = "",
read_buffer_size: USize val = 16384,
yield_after_reading: USize val = 16384,
yield_after_writing: USize val = 16384)
: TCPConnection tag^
Параметры
- auth: (AmbientAuth val | NetAuth val | TCPAuth val | TCPConnectAuth val)
- notify: TCPConnectionNotify iso
- host: String val
- service: String val
- from: String val = ""
- read_buffer_size: USize val = 16384
- yield_after_reading: USize val = 16384
- yield_after_writing: USize val = 16384
Возвращаемое значение
- TCPConnection tag^
Публичные методы
write
Запись одного набора байтов. Данные будут молча проигнорированы, если соединение ещё не установлено.
be write( data: (String val | Array[U8 val] val))
Параметры
writev
Запись последовательности последовательностей байтов. Данные будут молча проигнорированы, если соединение ещё не установлено.
be writev( data: ByteSeqIter val)
Параметры
- data: ByteSeqIter val
mute
Временно приостановить чтение из этого TCPConnection до тех пор, пока не будет вызван unmute.
be mute()
unmute
Начать чтение из этого TCPConnection снова после того, как он был заглушен.
be unmute()
set_notify
Изменить уведомляющее устройство.
be set_notify( notify: TCPConnectionNotify iso)
Параметры
- notify: TCPConnectionNotify iso
dispose
Закрыть соединение корректно после отправки всех записей.
be dispose()
Общедоступные функции
local_address
Возвращает локальный IP-адрес. Если это TCPConnection закрыто, возвращаемый адрес недействителен.
fun box local_address() : NetAddress val
Возвращаемое значение
- NetAddress val
remote_address
Возвращает удалённый IP-адрес. Если это TCPConnection закрыто, возвращаемый адрес недействителен.
fun box remote_address() : NetAddress val
Возвращаемое значение
- NetAddress val
expect
Вызов received на уведомляющем устройстве должен содержать ровно qty байтов. Если qty равно нулю, вызов может содержать любое количество данных. Это не имеет никакого эффекта, если вызывается в обратном вызове уведомляющего устройства sent.
Возникает ошибка, если qty превышает максимальный размер буфера, указанный в read_buffer_size при создании соединения.
fun ref expect( qty: USize val = 0) : None val ?
Параметры
- qty: USize val = 0
Возвращаемое значение
- None val ?
set_nodelay
Включить/выключить Nagle. По умолчанию включено. Это можно установить только на подключённом сокете.
fun ref set_nodelay( state: Bool val) : None val
Параметры
- state: Bool val
Возвращаемое значение
- None val
set_keepalive
Устанавливает таймаут TCP keepalive приблизительно до secs секунд. Точное время зависит от операционной системы. Если secs равно нулю, TCP keepalive отключен. TCP keepalive отключен по умолчанию. Это можно установить только на подключённом сокете.
fun ref set_keepalive( secs: U32 val) : None val
Параметры
- secs: U32 val
Возвращаемое значение
- None val
write_final
Напишите как можно больше в сокет. Установите _writeable в false в случае, если не всё было написано. В случае ошибки, закройте соединение. Это для данных, которые уже были преобразованы уведомлением. Данные будут молча отброшены, если соединение ещё не установлено.
fun ref write_final( data: (String val | Array[U8 val] val)) : None val
Параметры
Возвращаемое значение
- None val
close
Попытаться выполнить плавную остановку. Не принимать новые записи. Если соединение не заглушено, то мы не закончим закрытие, пока не получим чтение нулевой длины. Если соединение заглушено, выполните жёсткое закрытие и немедленно завершите работу.
fun ref close() : None val
Возвращаемое значение
- None val
hard_close
При возникновении ошибки выполните незапланированное закрытие.
fun ref hard_close() : None val
Возвращаемое значение
- None val
getsockopt
Общий оболочный интерфейс для TCP-сокетов к системному вызову getsockopt(2).
Вызывающий метод должен предоставить массив, предварительно выделенный размером не менее максимального размера структуры данных, которую ядро может вернуть для запрошенного параметра.
В случае успешного выполнения системного вызова эта функция возвращает пару: 1. Целое число 0. 2. Массив данных, возвращённых 4-м аргументом системного вызова. Его размер указан ядром через 5-й аргумент системного вызова.
В случае неудачи системного вызова эта функция возвращает пару: 1. Значение errno. 2. Неопределённое значение, которое следует игнорировать.
Пример использования:
// connected() is a callback function for class TCPConnectionNotify
fun ref connected(conn: TCPConnection ref) =>
match conn.getsockopt(OSSockOpt.sol_socket(), OSSockOpt.so_rcvbuf(), 4)
| (0, let gbytes: Array[U8] iso) =>
try
let br = Reader.create().>append(consume gbytes)
ifdef littleendian then
let buffer_size = br.u32_le()?
else
let buffer_size = br.u32_be()?
end
end
| (let errno: U32, _) =>
// System call failed
end
fun ref getsockopt( level: I32 val, option_name: I32 val, option_max_size: USize val = 4) : (U32 val , Array[U8 val] iso^)
Параметры
Возвращаемое значение
getsockopt_u32
Обёртка для TCP-сокетов к системному вызову getsockopt(2), где возвращаемое ядром значение параметра — тип C uint32_t,/ тип Pony U32.
В случае успешного выполнения системного вызова эта функция возвращает пару: 1. Целое число 0. 2. Значение *option_value, возвращённое ядром и преобразованное в тип Pony U32.
В случае неудачи системного вызова эта функция возвращает пару: 1. Значение errno. 2. Неопределённое значение, которое следует игнорировать.
fun ref getsockopt_u32( level: I32 val, option_name: I32 val) : (U32 val , U32 val)
Параметры
Возвращаемое значение
setsockopt
Общий оболочный интерфейс для TCP-сокетов к системному вызову setsockopt(2).
Вызывающий метод отвечает за правильный размер и содержимое байтов массива option для запрошенного level и option_name, включая использование соответствующего байтового порядка машины.
Эта функция возвращает 0 в случае успеха, в противном случае — значение errno в случае неудачи.
Пример использования:
// connected() is a callback function for class TCPConnectionNotify
fun ref connected(conn: TCPConnection ref) =>
let sb = Writer
sb.u32_le(7744) // Our desired socket buffer size
let sbytes = Array[U8]
for bs in sb.done().values() do
sbytes.append(bs)
end
match conn.setsockopt(OSSockOpt.sol_socket(), OSSockOpt.so_rcvbuf(), sbytes)
| 0 =>
// System call was successful
| let errno: U32 =>
// System call failed
end
fun ref setsockopt( level: I32 val, option_name: I32 val, option: Array[U8 val] ref) : U32 val
Параметры
Возвращаемое значение
- U32 val
setsockopt_u32
Общий оболочный интерфейс для TCP-сокетов к системному вызову setsockopt(2), где ядро ожидает значение параметра типа C uint32_t,/ типа Pony U32.
Эта функция возвращает 0 в случае успеха, в противном случае — значение errno в случае неудачи.
fun ref setsockopt_u32( level: I32 val, option_name: I32 val, option: U32 val) : U32 val
Параметры
Возвращаемое значение
- U32 val
get_so_error
Обёртка для системного вызова getsockopt(fd, SOL_SOCKET, SO_ERROR, ...)
fun ref get_so_error() : (U32 val , U32 val)
Возвращаемое значение
get_so_rcvbuf
Обёртка для системного вызова getsockopt(fd, SOL_SOCKET, SO_RCVBUF, ...)
fun ref get_so_rcvbuf() : (U32 val , U32 val)
Возвращаемое значение
get_so_sndbuf
Обёртка для системного вызова getsockopt(fd, SOL_SOCKET, SO_SNDBUF, ...)
fun ref get_so_sndbuf() : (U32 val , U32 val)
Возвращаемое значение
get_tcp_nodelay
Обёртка для системного вызова getsockopt(fd, SOL_SOCKET, TCP_NODELAY, ...)
fun ref get_tcp_nodelay() : (U32 val , U32 val)
Возвращаемое значение
set_so_rcvbuf
Обёртка для системного вызова setsockopt(fd, SOL_SOCKET, SO_RCVBUF, ...)
fun ref set_so_rcvbuf( bufsize: U32 val) : U32 val
Параметры
- bufsize: U32 val
Возвращаемое значение
- U32 val
set_so_sndbuf
Обёртка для системного вызова setsockopt(fd, SOL_SOCKET, SO_SNDBUF, ...)
fun ref set_so_sndbuf( bufsize: U32 val) : U32 val
Параметры
- bufsize: U32 val
Возвращаемое значение
- U32 val
set_tcp_nodelay
Обёртка для системного вызова setsockopt(fd, SOL_SOCKET, TCP_NODELAY, ...)
fun ref set_tcp_nodelay( state: Bool val) : U32 val
Параметры
- state: Bool val
Возвращаемое значение
- U32 val
© 2016-2020, The Pony Developers
© 2014-2015, Causality Ltd.
Licensed under the BSD 2-Clause License.
https://stdlib.ponylang.io/net-TCPConnection