Spec-Zone.ru › Python 3.13

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

Предисловие

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

Отправить байты данных удалённому узлу, заданному адресом (зависимый от транспорта целевой адрес). Если адрес равен None, данные отправляются на целевой адрес, заданный при создании транспорта.

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

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

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

Протоколы Datagram

Экземпляры протокола Datagram должны создаваться фабриками протоколов, переданными методу 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; 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–2024 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.13/library/asyncio-protocol.html

Spec-Zone.ru

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