Транспортные средства и Протоколы
Предисловие
Транспортные средства и Протоколы используются в низкоуровневых API цикла событий, таких как loop.create_connection(). Они используют стиль программирования на основе обратных вызовов и позволяют создавать высокопроизводительные реализации сетевых или межпроцессных протоколов (например, HTTP).
По существу, транспортные средства и протоколы должны использоваться только в библиотеках и фреймворках, а не в приложениях asyncio высокого уровня.
Эта страница документации охватывает как Транспортные средства, так и Протоколы.
Введение
На самом высоком уровне транспортное средство отвечает за то, как передаются байты, в то время как протокол определяет, какие байты передавать (и до некоторой степени, когда).
Другой способ сказать то же самое: транспортное средство является абстракцией сокета (или аналогичного конечного пункта ввода-вывода), а протокол является абстракцией приложения с точки зрения транспортного средства.
Еще одна точка зрения заключается в том, что интерфейсы транспортного средства и протокола вместе определяют абстрактный интерфейс для использования сетевого ввода-вывода и межпроцессного ввода-вывода.
Всегда существует однозначное соответствие между объектами транспортного средства и протокола: протокол вызывает методы транспортного средства для отправки данных, а транспортное средство вызывает методы протокола для передачи ему полученных данных.
Большинство методов цикла событий, ориентированных на установление соединения (таких как loop.create_connection()) обычно принимают аргумент protocol_factory, используемый для создания объекта Protocol для принятого соединения, представленного объектом Transport. Такие методы обычно возвращают кортеж (transport, protocol).
Содержание
Эта страница документации содержит следующие разделы:
- В разделе Транспортные средства документируются классы asyncio
BaseTransport,ReadTransport,WriteTransport,Transport,DatagramTransportиSubprocessTransport. - В разделе Протоколы документируются классы asyncio
BaseProtocol,Protocol,BufferedProtocol,DatagramProtocolиSubprocessProtocol. - Раздел Примеры демонстрирует, как работать с транспортными средствами, протоколами и API цикла событий низкого уровня.
Транспортные средства
Исходный код: Lib/asyncio/transports.py
Транспортные средства — это классы, предоставляемые модулем asyncio, для абстрагирования различных видов каналов связи.
Объекты транспортных средств всегда создаются циклом событий asyncio циклом событий.
asyncio реализует транспортные средства для TCP, UDP, SSL и каналов передачи данных процессам-потомкам. Доступные методы транспортного средства зависят от его типа.
Классы транспортных средств не являются потокобезопасными.
Иерархия транспортных средств
-
class asyncio.BaseTransport -
Базовый класс для всех транспортных средств. Содержит методы, общие для всех транспортных средств asyncio.
-
class asyncio.WriteTransport(BaseTransport) -
Базовый транспорт для односторонних подключений (только запись).
Экземпляры класса WriteTransport возвращаются методом цикла событий
loop.connect_write_pipe()и также используются в методах, связанных с дочерними процессами, таких какloop.subprocess_exec().
-
class asyncio.ReadTransport(BaseTransport) -
Базовый транспорт для односторонних подключений (только чтение).
Экземпляры класса ReadTransport возвращаются методом цикла событий
loop.connect_read_pipe()и также используются в методах, связанных с дочерними процессами, таких какloop.subprocess_exec().
-
class asyncio.Transport(WriteTransport, ReadTransport) -
Интерфейс, представляющий двусторонний транспорт, например, TCP-соединение.
Пользователь не создаёт экземпляр транспортного средства напрямую; он вызывает вспомогательную функцию, передавая ей фабрику протокола и другую информацию, необходимую для создания транспортного средства и протокола.
Экземпляры класса Transport возвращаются или используются методами цикла событий, такими как
loop.create_connection(),loop.create_unix_connection(),loop.create_server(),loop.sendfile()и др.
-
class asyncio.DatagramTransport(BaseTransport) -
Транспорт для соединений с использованием дейтаграмм (UDP).
Экземпляры класса DatagramTransport возвращаются методом цикла событий
loop.create_datagram_endpoint().
-
class asyncio.SubprocessTransport(BaseTransport) -
Абстракция для представления соединения между родительским и дочерним процессами ОС.
Экземпляры класса SubprocessTransport возвращаются методами цикла событий
loop.subprocess_shell()иloop.subprocess_exec().
Базовый транспорт
-
BaseTransport.close() -
Закрыть транспортное средство.
Если у транспортного средства есть буфер для исходящих данных, буферизованные данные будут асинхронно сброшены. Больше данных получено не будет. После сброса всех буферизованных данных будет вызван метод протокола
protocol.connection_lost()с аргументомNone. После закрытия транспортное средство использовать нельзя.
-
BaseTransport.is_closing() -
Возвращает
True, если транспорт закрывается или уже закрыт.
-
BaseTransport.get_extra_info(name, default=None) -
Возвращает информацию о транспортном средстве или связанных с ним ресурсах.
name — строка, представляющая запрашиваемую информацию, специфичную для транспортного средства.
default — значение, которое будет возвращено, если информация недоступна или если транспортное средство не поддерживает запросы с помощью данного стороннего цикла событий или на текущей платформе.
Например, следующий код пытается получить объект сокета, используемый транспортным средством:
sock = transport.get_extra_info('socket') if sock is not None: print(sock.getsockopt(...))Категории информации, которые могут быть запрошены для некоторых транспортных средств:
-
сокет:
-
'peername': удалённый адрес, к которому подключён сокет, результат вызоваsocket.socket.getpeername()(Noneпри ошибке) -
'socket': экземплярsocket.socket -
'sockname': собственный адрес сокета, результат вызоваsocket.socket.getsockname()
-
-
сокет SSL:
-
'compression': алгоритм сжатия в виде строки илиNoneесли сжатие не используется; результат вызоваssl.SSLSocket.compression() -
'cipher': кортеж из трёх значений, содержащий имя используемого шифра, версию протокола SSL, определяющую его использование, и число битов секрета; результат вызоваssl.SSLSocket.cipher() -
'peercert': сертификат удалённого узла; результат вызоваssl.SSLSocket.getpeercert() -
'sslcontext': экземплярssl.SSLContext -
'ssl_object': экземплярssl.SSLObjectилиssl.SSLSocket
-
-
пайп:
-
'pipe': объект пайпа
-
-
подпроцесс:
-
'subprocess': экземплярsubprocess.Popen
-
-
-
BaseTransport.set_protocol(protocol) -
Установить новый протокол.
Переключение протоколов должно осуществляться только в тех случаях, когда оба протокола документированно поддерживают такое переключение.
-
BaseTransport.get_protocol() -
Возвращает текущий протокол.
Транспортные средства только для чтения
-
ReadTransport.is_reading() -
Возвращает
True, если транспорт получает новые данные.Добавлен в версии 3.7.
-
ReadTransport.pause_reading() -
Приостановить приём данных транспортным средством. Данные не будут передаваться методу протокола
protocol.data_received()до тех пор, пока не будет вызван методresume_reading().Изменено в версии 3.7: Метод является идемпотентным, то есть его можно вызывать, когда транспорт уже приостановлен или закрыт.
-
ReadTransport.resume_reading() -
Возобновить приём данных. Метод протокола
protocol.data_received()будет снова вызван, если доступны данные для чтения.Изменено в версии 3.7: Метод является идемпотентным, то есть его можно вызывать, когда транспорт уже читает.
Только для записи транспорты
-
WriteTransport.abort() -
Немедленно закрыть транспорт, не дожидаясь завершения ожидающих операций. Буферизованные данные будут потеряны. Больше данных не будет получено. Метод протокола
protocol.connection_lost()в конечном итоге будет вызван сNoneв качестве аргумента.
-
WriteTransport.can_write_eof() -
Возвращает
True, если транспорт поддерживаетwrite_eof(), иFalseв противном случае.
-
WriteTransport.get_write_buffer_size() -
Возвращает текущий размер буфера вывода, используемого транспортом.
-
WriteTransport.get_write_buffer_limits() -
Получить верхний и нижний пороги для управления потоком записи. Возвращает кортеж
(low, high), где нижний и верхний — положительные числа байтов.Используйте
set_write_buffer_limits()для установки ограничений.Введено в версии 3.4.2.
-
WriteTransport.set_write_buffer_limits(high=None, low=None) -
Установить верхний и нижний пороги для управления потоком записи.
Эти два значения (измеряемые в количестве байтов) управляют тем, когда вызываются методы протокола
protocol.pause_writing()иprotocol.resume_writing(). Если указаны, нижний порог должен быть меньше или равен верхнему. Ни верхний, ни нижний порог не могут быть отрицательными.Вызов
pause_writing()происходит, когда размер буфера становится больше или равен верхнему значению. Если запись была приостановлена, вызовresume_writing()происходит, когда размер буфера становится меньше или равен нижнему значению.Значения по умолчанию зависят от реализации. Если задан только верхний порог, нижний порог устанавливается по умолчанию до значения, зависящего от реализации, которое меньше или равно верхнему порогу. Установка верхнего порога в ноль также устанавливает нижний порог в ноль и вызывает
pause_writing()всякий раз, когда буфер становится непустым. Установка нижнего порога в ноль вызываетresume_writing()только после того, как буфер станет пустым. Использование нулевых значений для любого из порогов обычно не является оптимальным, поскольку это уменьшает возможности одновременного выполнения операций ввода/вывода и вычислений.Получить ограничения можно с помощью
get_write_buffer_limits().
-
WriteTransport.write(data) -
Записать несколько байтов data в транспорт.
Этот метод не блокирует; он буферизует данные и организует их асинхронную отправку.
-
WriteTransport.writelines(list_of_data) -
Записать список (или любой итерируемый объект) байтов данных в транспорт. Это функционально эквивалентно вызову
write()для каждого элемента, возвращаемого итерируемым объектом, но может быть реализовано более эффективно.
-
WriteTransport.write_eof() -
Закрыть конец записи транспорта после сброса всех буферизованных данных. Прием данных все еще возможен.
Этот метод может вызвать
NotImplementedError, если транспорт (например, SSL) не поддерживает полузакрытые соединения.
Транспорты для дейтаграмм
-
DatagramTransport.sendto(data, addr=None) -
Отправить байты data удалённому узлу, указанному в addr (адрес, зависящий от транспорта). Если addr равен
None, данные отправляются по адресу, заданному при создании транспорта.Этот метод не блокирует; он буферизует данные и организует их асинхронную отправку.
-
DatagramTransport.abort() -
Немедленно закрыть транспорт, не дожидаясь завершения ожидающих операций. Буферизованные данные будут потеряны. Больше данных не будет получено. Метод протокола
protocol.connection_lost()в конечном итоге будет вызван сNoneв качестве аргумента.
Транспорты для дочерних процессов
-
SubprocessTransport.get_pid() -
Возвращает идентификатор процесса дочернего процесса в виде целого числа.
-
SubprocessTransport.get_pipe_transport(fd) -
Возвращает транспорт для канала связи, соответствующий целочисленному дескриптору файла fd:
-
0: транспорт потокового чтения стандартного ввода (stdin) илиNone, если дочерний процесс не был создан сstdin=PIPE -
1: транспорт потоковой записи стандартного вывода (stdout) илиNone, если дочерний процесс не был создан сstdout=PIPE -
2: транспорт потоковой записи стандартной ошибки (stderr) илиNone, если дочерний процесс не был создан сstderr=PIPE - другой fd:
None
-
-
SubprocessTransport.get_returncode() -
Возвращает код возврата дочернего процесса в виде целого числа или
None, если он ещё не возвращён, что аналогично атрибутуsubprocess.Popen.returncode.
-
SubprocessTransport.kill() -
Убить дочерний процесс.
В системах POSIX функция отправляет SIGKILL дочернему процессу. В Windows этот метод является псевдонимом для
terminate().См. также
subprocess.Popen.kill().
-
SubprocessTransport.send_signal(signal) -
Отправить дочернему процессу номер сигнала, как в
subprocess.Popen.send_signal().
-
SubprocessTransport.terminate() -
Остановить дочерний процесс.
В системах POSIX этот метод отправляет SIGTERM дочернему процессу. В Windows вызывается функция Windows API TerminateProcess() для остановки дочернего процесса.
См. также
subprocess.Popen.terminate().
-
SubprocessTransport.close() -
Убить дочерний процесс, вызвав метод
kill().Если дочерний процесс ещё не завершился, и закрыть транспорты для stdin, stdout и stderr каналов.
Протоколы
Исходный код: Lib/asyncio/protocols.py
asyncio предоставляет набор абстрактных базовых классов, которые следует использовать для реализации сетевых протоколов. Эти классы предназначены для совместного использования с транспортами.
Подклассы абстрактных базовых классов протоколов могут реализовывать некоторые или все методы. Все эти методы являются обратными вызовами: они вызываются транспортом при определенных событиях, например, при получении данных. Базовый метод протокола должен вызываться соответствующим транспортом.
Базовые протоколы
-
class asyncio.BaseProtocol -
Базовый протокол с методами, общими для всех протоколов.
-
class asyncio.Protocol(BaseProtocol) -
Базовый класс для реализации потоковых протоколов (TCP, сокеты Unix и т. д.).
-
class asyncio.BufferedProtocol(BaseProtocol) -
Базовый класс для реализации потоковых протоколов с ручным управлением буфером приема.
-
class asyncio.DatagramProtocol(BaseProtocol) -
Базовый класс для реализации протоколов дейтаграмм (UDP).
-
class asyncio.SubprocessProtocol(BaseProtocol) -
Базовый класс для реализации протоколов, взаимодействующих с дочерними процессами (односторонние каналы).
Базовый протокол
Все протоколы asyncio могут реализовывать обратные вызовы базового протокола.
Обратные вызовы подключения
Обратные вызовы подключения вызываются для всех протоколов ровно один раз при успешном подключении. Все остальные обратные вызовы протокола могут вызываться только между этими двумя методами.
-
BaseProtocol.connection_made(transport) -
Вызывается при установлении соединения.
Аргумент transport — это транспорт, представляющий соединение. Протокол отвечает за сохранение ссылки на свой транспорт.
-
BaseProtocol.connection_lost(exc) -
Вызывается при потере или закрытии соединения.
Аргументом является либо объект исключения, либо
None. Последнее означает, что получен обычный EOF или соединение было прервано или закрыто этой стороной соединения.
Обратные вызовы управления потоком
Обратные вызовы управления потоком могут вызываться транспортом для приостановки или возобновления записи, выполняемой протоколом.
Для получения дополнительной информации см. документацию метода set_write_buffer_limits().
-
BaseProtocol.pause_writing() -
Вызывается, когда буфер транспорта превышает верхний порог.
-
BaseProtocol.resume_writing() -
Вызывается, когда буфер транспорта опускается ниже нижнего порога.
Если размер буфера равен верхнему порогу, pause_writing() не вызывается: размер буфера должен строго превышать порог.
И наоборот, resume_writing() вызывается, когда размер буфера равен или меньше нижнего порога. Эти конечные условия важны для обеспечения корректной работы, когда один из порогов равен нулю.
Потоковые протоколы
Методы событий, такие как loop.create_server(), loop.create_unix_server(), loop.create_connection(), loop.create_unix_connection(), loop.connect_accepted_socket(), loop.connect_read_pipe() и loop.connect_write_pipe() принимают фабрики, которые возвращают потоковые протоколы.
-
Protocol.data_received(data) -
Вызывается при получении данных. data — это непустой объект bytes, содержащий входящие данные.
Буферизация, фрагментация или повторная сборка данных зависят от транспорта. В общем случае не следует полагаться на определенные семантики, а вместо этого следует делать ваш анализ общим и гибким. Однако данные всегда принимаются в правильном порядке.
Метод может вызываться любое количество раз во время открытого соединения.
Однако
protocol.eof_received()вызывается не более одного раза. После вызоваeof_received()вызовdata_received()больше не происходит.
-
Protocol.eof_received() -
Вызывается, когда другой конец сигнализирует, что больше не будет отправлять данные (например, вызвав
transport.write_eof(), если другой конец также использует asyncio).Этот метод может вернуть значение false (включая
None), в этом случае транспорт закроет себя. И наоборот, если этот метод возвращает true, протокол определяет, закрыть ли транспорт. Поскольку реализация по умолчанию возвращаетNone, она неявно закрывает соединение.Некоторые транспорты, включая SSL, не поддерживают полузакрытые соединения, поэтому возвращение true из этого метода приведет к закрытию соединения.
Схема состояний:
start -> connection_made
[-> data_received]*
[-> eof_received]?
-> connection_lost -> end
Буферизованные потоковые протоколы
Введено в версии 3.7.
Буферизованные протоколы могут использоваться с любым методом цикла событий, поддерживающим Потоковые протоколы.
Реализации BufferedProtocol позволяют явно выделять и управлять буфером приема. Циклы событий могут использовать буфер, предоставленный протоколом, для предотвращения ненужных копий данных. Это может привести к заметному улучшению производительности для протоколов, принимающих большие объемы данных. Сложные реализации протоколов могут значительно сократить число выделений буферов.
Следующие обратные вызовы вызываются на экземплярах BufferedProtocol:
-
BufferedProtocol.get_buffer(sizehint) -
Вызывается для выделения нового буфера приема.
sizehint — это рекомендуемый минимальный размер возвращаемого буфера. Приемлемо возвращать буферы меньшего или большего размера, чем предполагает sizehint. Если установлено в -1, размер буфера может быть произвольным. Ошибка — возвращение буфера с нулевым размером.
get_buffer()должен возвращать объект, реализующий протокол буфера.
-
BufferedProtocol.buffer_updated(nbytes) -
Вызывается, когда буфер был обновлен полученными данными.
nbytes — общее количество байтов, записанных в буфер.
-
BufferedProtocol.eof_received() -
См. документацию метода
protocol.eof_received().
get_buffer() может вызываться любое количество раз во время соединения. Однако protocol.eof_received() вызывается не более одного раза и, если вызывается, get_buffer() и buffer_updated() больше не будут вызываться после этого.
Схема состояний:
start -> connection_made
[-> get_buffer
[-> buffer_updated]?
]*
[-> eof_received]?
-> connection_lost -> end
Протоколы дейтаграмм
Экземпляры протокола дейтаграмм должны создаваться фабриками протоколов, переданными методу loop.create_datagram_endpoint().
-
DatagramProtocol.datagram_received(data, addr) -
Вызывается при получении дейтаграммы. data — это объект bytes, содержащий входящие данные. addr — адрес отправителя данных; точный формат зависит от транспорта.
-
DatagramProtocol.error_received(exc) -
Вызывается, когда предыдущая операция отправки или получения данных вызывает
OSError. exc — экземплярOSError.Этот метод вызывается в редких случаях, когда транспорт (например, UDP) обнаруживает, что дейтаграмму невозможно доставить получателю. Однако во многих случаях недоставляемые дейтаграммы будут просто проигнорированы.
Примечание
В системах BSD (macOS, FreeBSD и т. д.) управление потоком для протоколов дейтаграмм не поддерживается, поскольку нет надежного способа обнаружить сбои при отправке, вызванные отправкой слишком большого количества пакетов.
Сокет всегда выглядит «готовым», и лишние пакеты отбрасываются. OSError с errno, установленным в значение errno.ENOBUFS, может быть или не быть вызван; если он вызван, он будет передан в DatagramProtocol.error_received(), в противном случае будет проигнорирован.
Протоколы подпроцессов
Экземпляры протоколов подпроцессов должны создаваться фабриками протоколов, переданными методам loop.subprocess_exec() и loop.subprocess_shell().
-
SubprocessProtocol.pipe_data_received(fd, data) -
Вызывается, когда дочерний процесс записывает данные в свою трубу stdout или stderr.
fd — целое число, описывающее дескриптор файла трубы.
data — непустой объект bytes, содержащий полученные данные.
-
SubprocessProtocol.pipe_connection_lost(fd, exc) -
Вызывается, когда одна из труб, общающихся с дочерним процессом, закрывается.
fd — целое число, описывающее закрытый дескриптор файла.
-
SubprocessProtocol.process_exited() -
Вызывается, когда дочерний процесс завершил работу.
Он может быть вызван до методов
pipe_data_received()иpipe_connection_lost().
Примеры
Эхо-сервер TCP
Создайте эхо-сервер TCP, используя метод loop.create_server(), отправьте обратно полученные данные и закройте соединение:
import asyncio
class EchoServerProtocol(asyncio.Protocol):
def connection_made(self, transport):
peername = transport.get_extra_info('peername')
print('Connection from {}'.format(peername))
self.transport = transport
def data_received(self, data):
message = data.decode()
print('Data received: {!r}'.format(message))
print('Send: {!r}'.format(message))
self.transport.write(data)
print('Close the client socket')
self.transport.close()
async def main():
# Get a reference to the event loop as we plan to use
# low-level APIs.
loop = asyncio.get_running_loop()
server = await loop.create_server(
lambda: EchoServerProtocol(),
'127.0.0.1', 8888)
async with server:
await server.serve_forever()
asyncio.run(main())
См. также
В примере эхо-сервер TCP с использованием потоков используется высокоуровневая функция asyncio.start_server().
Эхо-клиент TCP
Эхо-клиент TCP, используя метод loop.create_connection(), отправляет данные и ожидает, пока соединение не закроется:
import asyncio
class EchoClientProtocol(asyncio.Protocol):
def __init__(self, message, on_con_lost):
self.message = message
self.on_con_lost = on_con_lost
def connection_made(self, transport):
transport.write(self.message.encode())
print('Data sent: {!r}'.format(self.message))
def data_received(self, data):
print('Data received: {!r}'.format(data.decode()))
def connection_lost(self, exc):
print('The server closed the connection')
self.on_con_lost.set_result(True)
async def main():
# Get a reference to the event loop as we plan to use
# low-level APIs.
loop = asyncio.get_running_loop()
on_con_lost = loop.create_future()
message = 'Hello World!'
transport, protocol = await loop.create_connection(
lambda: EchoClientProtocol(message, on_con_lost),
'127.0.0.1', 8888)
# Wait until the protocol signals that the connection
# is lost and close the transport.
try:
await on_con_lost
finally:
transport.close()
asyncio.run(main())
См. также
В примере эхо-клиент TCP с использованием потоков используется высокоуровневая функция asyncio.open_connection().
Эхо-сервер UDP
Эхо-сервер UDP, использующий метод loop.create_datagram_endpoint(), отправляет обратно полученные данные:
import asyncio
class EchoServerProtocol:
def connection_made(self, transport):
self.transport = transport
def datagram_received(self, data, addr):
message = data.decode()
print('Received %r from %s' % (message, addr))
print('Send %r to %s' % (message, addr))
self.transport.sendto(data, addr)
async def main():
print("Starting UDP server")
# Get a reference to the event loop as we plan to use
# low-level APIs.
loop = asyncio.get_running_loop()
# One protocol instance will be created to serve all
# client requests.
transport, protocol = await loop.create_datagram_endpoint(
lambda: EchoServerProtocol(),
local_addr=('127.0.0.1', 9999))
try:
await asyncio.sleep(3600) # Serve for 1 hour.
finally:
transport.close()
asyncio.run(main())
Эхо-клиент UDP
Эхо-клиент UDP, использующий метод loop.create_datagram_endpoint(), отправляет данные и закрывает транспорт при получении ответа:
import asyncio
class EchoClientProtocol:
def __init__(self, message, on_con_lost):
self.message = message
self.on_con_lost = on_con_lost
self.transport = None
def connection_made(self, transport):
self.transport = transport
print('Send:', self.message)
self.transport.sendto(self.message.encode())
def datagram_received(self, data, addr):
print("Received:", data.decode())
print("Close the socket")
self.transport.close()
def error_received(self, exc):
print('Error received:', exc)
def connection_lost(self, exc):
print("Connection closed")
self.on_con_lost.set_result(True)
async def main():
# Get a reference to the event loop as we plan to use
# low-level APIs.
loop = asyncio.get_running_loop()
on_con_lost = loop.create_future()
message = "Hello World!"
transport, protocol = await loop.create_datagram_endpoint(
lambda: EchoClientProtocol(message, on_con_lost),
remote_addr=('127.0.0.1', 9999))
try:
await on_con_lost
finally:
transport.close()
asyncio.run(main())
Подключение к существующим сокетам
Ожидание, пока сокет получит данные, используя метод loop.create_connection() с протоколом:
import asyncio
import socket
class MyProtocol(asyncio.Protocol):
def __init__(self, on_con_lost):
self.transport = None
self.on_con_lost = on_con_lost
def connection_made(self, transport):
self.transport = transport
def data_received(self, data):
print("Received:", data.decode())
# We are done: close the transport;
# connection_lost() will be called automatically.
self.transport.close()
def connection_lost(self, exc):
# The socket has been closed
self.on_con_lost.set_result(True)
async def main():
# Get a reference to the event loop as we plan to use
# low-level APIs.
loop = asyncio.get_running_loop()
on_con_lost = loop.create_future()
# Create a pair of connected sockets
rsock, wsock = socket.socketpair()
# Register the socket to wait for data.
transport, protocol = await loop.create_connection(
lambda: MyProtocol(on_con_lost), sock=rsock)
# Simulate the reception of data from the network.
loop.call_soon(wsock.send, 'abc'.encode())
try:
await protocol.on_con_lost
finally:
transport.close()
wsock.close()
asyncio.run(main())
См. также
Пример мониторинг дескриптора файла на наличие событий чтения использует низкоуровневый метод loop.add_reader() для регистрации FD.
Пример регистрация открытого сокета для ожидания данных с помощью потоков использует высокоуровневые потоки, созданные функцией open_connection() в корутине.
loop.subprocess_exec() и SubprocessProtocol
Пример использования протокола подпроцесса для получения вывода подпроцесса и ожидания завершения работы подпроцесса.
Подпроцесс создается методом loop.subprocess_exec():
import asyncio
import sys
class DateProtocol(asyncio.SubprocessProtocol):
def __init__(self, exit_future):
self.exit_future = exit_future
self.output = bytearray()
self.pipe_closed = False
self.exited = False
def pipe_connection_lost(self, fd, exc):
self.pipe_closed = True
self.check_for_exit()
def pipe_data_received(self, fd, data):
self.output.extend(data)
def process_exited(self):
self.exited = True
# process_exited() method can be called before
# pipe_connection_lost() method: wait until both methods are
# called.
self.check_for_exit()
def check_for_exit(self):
if self.pipe_closed and self.exited:
self.exit_future.set_result(True)
async def get_date():
# Get a reference to the event loop as we plan to use
# low-level APIs.
loop = asyncio.get_running_loop()
code = 'import datetime; print(datetime.datetime.now())'
exit_future = asyncio.Future(loop=loop)
# Create the subprocess controlled by DateProtocol;
# redirect the standard output into a pipe.
transport, protocol = await loop.subprocess_exec(
lambda: DateProtocol(exit_future),
sys.executable, '-c', code,
stdin=None, stderr=None)
# Wait for the subprocess exit using the process_exited()
# method of the protocol.
await exit_future
# Close the stdout pipe.
transport.close()
# Read the output which was collected by the
# pipe_data_received() method of the protocol.
data = bytes(protocol.output)
return data.decode('ascii').rstrip()
date = asyncio.run(get_date())
print(f"Current date: {date}")
См. также аналогичный пример, написанный с использованием высокоуровневых API.
© 2001–2023 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.11/library/asyncio-protocol.html