Spec-Zone.ru › Python 3.14

Транспорты и протоколы

Предисловие

Транспорты и протоколы используются API цикла событий низкого уровня, например loop.create_connection(). В них применяется стиль программирования на основе обратных вызовов, что позволяет создавать высокопроизводительные реализации сетевых протоколов и протоколов IPC (например, HTTP).

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

На этой странице документации рассматриваются как транспорты, так и протоколы.

Введение

На самом общем уровне транспорт отвечает за то, как передаются байты, а протокол определяет, какие байты передавать (и в некоторой степени когда).

То же самое можно выразить иначе: транспорт — это абстракция сокета (или аналогичной конечной точки ввода-вывода), а протокол — абстракция приложения с точки зрения транспорта.

Ещё один взгляд: интерфейсы транспорта и протокола вместе определяют абстрактный интерфейс для сетевого и межпроцессного ввода-вывода.

Между объектами транспорта и протокола всегда существует связь 1:1: протокол вызывает методы транспорта для отправки данных, а транспорт вызывает методы протокола, чтобы передать ему полученные данные.

Большинство методов цикла событий, ориентированных на соединение (например, 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), где low и high — положительные значения в байтах.

Для установки пороговых уровней используйте set_write_buffer_limits().

Добавлено в версии 3.4.2.

WriteTransport.set_write_buffer_limits(high=None, low=None)

Установить верхний и нижний пороговые уровни для управления потоком записи.

Эти два значения (измеряемые в байтах) определяют, когда вызываются методы протокола protocol.pause_writing() и protocol.resume_writing(). Если указано нижнее пороговое значение, оно должно быть меньше или равно верхнему. Значения high и low не могут быть отрицательными.

Метод pause_writing() вызывается, когда размер буфера становится больше или равен значению high. Если запись приостановлена, метод resume_writing() вызывается, когда размер буфера становится меньше или равен значению low.

Значения по умолчанию зависят от реализации. Если указано только верхнее пороговое значение, нижнее по умолчанию устанавливается в зависящее от реализации значение, меньшее или равное верхнему. Если установить high равным нулю, значение low также станет нулевым, а метод pause_writing() будет вызываться каждый раз, когда буфер становится непустым. Если установить low равным нулю, метод 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, данные отправляются по адресу назначения, заданному при создании транспорта.

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

Изменено в версии 3.13: Этот метод можно вызывать с пустым объектом bytes, чтобы отправить датаграмму нулевой длины. Кроме того, расчёт размера буфера для управления потоком теперь учитывает заголовок датаграммы.

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)

Отправить подпроцессу сигнал с номером signal, как в subprocess.Popen.send_signal().

SubprocessTransport.terminate()

Остановить подпроцесс.

В системах POSIX этот метод отправляет подпроцессу сигнал SIGTERM. В Windows для остановки подпроцесса вызывается функция API Windows 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).

Этот метод может вернуть ложное значение (включая None), в этом случае транспорт закроется. Если же метод возвращает истинное значение, то решение о закрытии транспорта зависит от используемого протокола. Поскольку реализация по умолчанию возвращает None, соединение неявно закрывается.

Некоторые транспорты, включая SSL, не поддерживают полузакрытые соединения; в таком случае возврат истинного значения этим методом приведёт к закрытию соединения.

Конечный автомат:

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(
        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(
        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 as dt; print(dt.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 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.14/library/asyncio-protocol.html

Spec-Zone.ru

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