Spec-Zone.ru › Python 3.8

Транспортные средства и Протоколы

Предисловие

Транспортные средства и Протоколы используются в 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)

Записать некоторые байты данных в транспорт.

Этот метод не блокирует; он буферизует данные и организует их асинхронную отправку.

WriteTransport.writelines(list_of_data)

Записать список (или любой итерируемый объект) байтов данных в транспорт. Это функционально эквивалентно вызову write() для каждого элемента, возвращаемого итерируемым объектом, но может быть реализовано более эффективно.

WriteTransport.write_eof()

Закрыть конец записи транспорта после сброса всех буферизованных данных. Данные могут все еще приниматься.

Этот метод может вызывать исключение NotImplementedError, если транспорт (например, SSL) не поддерживает полузакрытые соединения.

Транспортные средства для дейтаграмм

DatagramTransport.sendto(data, addr=None)

Отправить байты данных удалённому узлу, заданному 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 могут реализовывать обратные вызовы Base Protocol.

Обратные вызовы подключения

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

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: Важно: это было добавлено в asyncio в Python 3.7 в предварительном порядке! Это экспериментальный API, который может быть изменён или удалён полностью в Python 3.8.

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

Реализации 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()

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

Примеры

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()

    def pipe_data_received(self, fd, data):
        self.output.extend(data)

    def process_exited(self):
        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–2022 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.8/library/asyncio-protocol.html

Spec-Zone.ru

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