Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -75,7 +75,7 @@ jobs:
strategy:
matrix:
os: [ubuntu-latest, macos-latest]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*"]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*", "cp311-*"]
cibw_arch: ["x86_64", "aarch64", "universal2"]
exclude:
- os: ubuntu-latest
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/tests.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -14,7 +14,7 @@ jobs:
runs-on: ${{ matrix.os }}
strategy:
matrix:
python-version: ["3.7", "3.8", "3.9", "3.10"]
python-version: ["3.7", "3.8", "3.9", "3.10", "3.11.0-rc.1"]
os: [ubuntu-latest, macos-latest]

env:
Expand Down
6 changes: 4 additions & 2 deletions setup.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,14 +21,16 @@
from setuptools.command.sdist import sdist


CYTHON_DEPENDENCY = 'Cython(>=0.29.24,<0.30.0)'
CYTHON_DEPENDENCY = 'Cython(>=0.29.32,<0.30.0)'

# Minimal dependencies required to test uvloop.
TEST_DEPENDENCIES = [
# pycodestyle is a dependency of flake8, but it must be frozen because
# their combination breaks too often
# (example breakage: https://gitlab.com/pycqa/flake8/issues/427)
'aiohttp',
# aiohttp doesn't support 3.11 yet,
# see https://github.com/aio-libs/aiohttp/issues/6600
'aiohttp ; python_version < "3.11"',
'flake8~=3.9.2',
'psutil',
'pycodestyle~=2.7.0',
Expand Down
8 changes: 6 additions & 2 deletions tests/test_base.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -540,7 +540,9 @@ class MyTask(asyncio.Task):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
self.loop.set_task_factory(factory)
Expand DownExpand Up@@ -577,7 +579,9 @@ def get_name(self):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
task = self.loop.create_task(coro(), name="mytask")
Expand Down
4 changes: 4 additions & 0 deletions tests/test_context.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -452,6 +452,10 @@ def close():
self._run_server_test(test, async_sock=True)

def test_create_ssl_server_manual_connection_lost(self):
if self.implementation == 'asyncio' and sys.version_info >= (3, 11, 0):
# TODO(fantix): fix for 3.11
raise unittest.SkipTest('should pass on 3.11')

async def test(proto, cvar, ssl_sock, **_):
def close():
cvar.set('closing')
Expand Down
37 changes: 0 additions & 37 deletions tests/test_tcp.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -652,43 +652,6 @@ async def runner():
self.assertIsNone(
self.loop.run_until_complete(connection_lost_called))

def test_context_run_segfault(self):
is_new = False
done = self.loop.create_future()

def server(sock):
sock.sendall(b'hello')

class Protocol(asyncio.Protocol):
def __init__(self):
self.transport = None

def connection_made(self, transport):
self.transport = transport

def data_received(self, data):
try:
self = weakref.ref(self)
nonlocal is_new
if is_new:
done.set_result(data)
else:
is_new = True
new_proto = Protocol()
self().transport.set_protocol(new_proto)
new_proto.connection_made(self().transport)
new_proto.data_received(data)
except Exception as e:
done.set_exception(e)

async def test(addr):
await self.loop.create_connection(Protocol, *addr)
data = await done
self.assertEqual(data, b'hello')

with self.tcp_server(server) as srv:
self.loop.run_until_complete(test(srv.addr))


class Test_UV_TCP(_TestTCP, tb.UVTestCase):

Expand Down
4 changes: 4 additions & 0 deletions uvloop/loop.pyi
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import asyncio
import ssl
import sys
from socket import AddressFamily, SocketKind, _Address, _RetAddress, socket
from typing import (
IO,
Expand DownExpand Up@@ -210,6 +211,9 @@ class Loop:
async def sock_sendall(self, sock: socket, data: bytes) -> None: ...
async def sock_accept(self, sock: socket) -> Tuple[socket, _RetAddress]: ...
async def sock_connect(self, sock: socket, address: _Address) -> None: ...
async def sock_recvfrom(self, sock: socket, bufsize: int) -> bytes: ...
async def sock_recvfrom_into(self, sock: socket, buf: bytearray, nbytes: int = ...) -> int: ...
async def sock_sendto(self, sock: socket, data: bytes, address: _Address) -> None: ...
async def connect_accepted_socket(
self,
protocol_factory: Callable[[], _ProtocolT],
Expand Down
37 changes: 33 additions & 4 deletions uvloop/loop.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@ include "errors.pyx"

cdef:
int PY39 = PY_VERSION_HEX >= 0x03090000
int PY311 = PY_VERSION_HEX >= 0x030b0000
uint64_t MAX_SLEEP = 3600 * 24 * 365 * 100


Expand DownExpand Up@@ -1413,19 +1414,35 @@ cdef class Loop:
"""Create a Future object attached to the loop."""
return self._new_future()

def create_task(self, coro, *, name=None):
def create_task(self, coro, *, name=None, context=None):
"""Schedule a coroutine object.

Return a task object.

If name is not None, task.set_name(name) will be called if the task
object has the set_name attribute, true for default Task in Python 3.8.

An optional keyword-only context argument allows specifying a custom
contextvars.Context for the coro to run in. The current context copy is
created when no context is provided.
"""
self._check_closed()
if self._task_factory is None:
task = aio_Task(coro, loop=self)
if PY311:
if self._task_factory is None:
task = aio_Task(coro, loop=self, context=context)
else:
task = self._task_factory(self, coro, context=context)
else:
task = self._task_factory(self, coro)
if context is None:
if self._task_factory is None:
task = aio_Task(coro, loop=self)
else:
task = self._task_factory(self, coro)
else:
if self._task_factory is None:
task = context.run(aio_Task, coro, self)
else:
task = context.run(self._task_factory, self, coro)

# copied from asyncio.tasks._set_task_name (bpo-34270)
if name is not None:
Expand DownExpand Up@@ -2604,6 +2621,18 @@ cdef class Loop:
finally:
socket_dec_io_ref(sock)

@cython.iterable_coroutine
async def sock_recvfrom(self, sock, bufsize):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_recvfrom_into(self, sock, buf, nbytes=0):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_sendto(self, sock, data, address):
raise NotImplementedError

@cython.iterable_coroutine
async def connect_accepted_socket(self, protocol_factory, sock, *,
ssl=None,
Expand Down
4 changes: 2 additions & 2 deletions uvloop/pseudosock.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,8 +41,8 @@ cdef class PseudoSocket:

def __repr__(self):
s = ("<uvloop.PseudoSocket fd={}, family={!s}, "
"type={!s}, proto={}").format(self.fileno(), self.family,
self.type, self.proto)
"type={!s}, proto={}").format(self.fileno(), self.family.name,
self.type.name, self.proto)

if self._fd != -1:
try:
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all \u003cpre\u003e\u003ccode\u003e blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks"); } } catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); } })(); (function(){ try { var __m = "github.com"; var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -75,7 +75,7 @@ jobs:
strategy:
matrix:
os: [ubuntu-latest, macos-latest]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*"]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*", "cp311-*"]
cibw_arch: ["x86_64", "aarch64", "universal2"]
exclude:
- os: ubuntu-latest
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/tests.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -14,7 +14,7 @@ jobs:
runs-on: ${{ matrix.os }}
strategy:
matrix:
python-version: ["3.7", "3.8", "3.9", "3.10"]
python-version: ["3.7", "3.8", "3.9", "3.10", "3.11.0-rc.1"]
os: [ubuntu-latest, macos-latest]

env:
Expand Down
6 changes: 4 additions & 2 deletions setup.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,14 +21,16 @@
from setuptools.command.sdist import sdist


CYTHON_DEPENDENCY = 'Cython(>=0.29.24,<0.30.0)'
CYTHON_DEPENDENCY = 'Cython(>=0.29.32,<0.30.0)'

# Minimal dependencies required to test uvloop.
TEST_DEPENDENCIES = [
# pycodestyle is a dependency of flake8, but it must be frozen because
# their combination breaks too often
# (example breakage: https://gitlab.com/pycqa/flake8/issues/427)
'aiohttp',
# aiohttp doesn't support 3.11 yet,
# see https://github.com/aio-libs/aiohttp/issues/6600
'aiohttp ; python_version < "3.11"',
'flake8~=3.9.2',
'psutil',
'pycodestyle~=2.7.0',
Expand Down
8 changes: 6 additions & 2 deletions tests/test_base.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -540,7 +540,9 @@ class MyTask(asyncio.Task):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
self.loop.set_task_factory(factory)
Expand DownExpand Up@@ -577,7 +579,9 @@ def get_name(self):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
task = self.loop.create_task(coro(), name="mytask")
Expand Down
4 changes: 4 additions & 0 deletions tests/test_context.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -452,6 +452,10 @@ def close():
self._run_server_test(test, async_sock=True)

def test_create_ssl_server_manual_connection_lost(self):
if self.implementation == 'asyncio' and sys.version_info >= (3, 11, 0):
# TODO(fantix): fix for 3.11
raise unittest.SkipTest('should pass on 3.11')

async def test(proto, cvar, ssl_sock, **_):
def close():
cvar.set('closing')
Expand Down
37 changes: 0 additions & 37 deletions tests/test_tcp.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -652,43 +652,6 @@ async def runner():
self.assertIsNone(
self.loop.run_until_complete(connection_lost_called))

def test_context_run_segfault(self):
is_new = False
done = self.loop.create_future()

def server(sock):
sock.sendall(b'hello')

class Protocol(asyncio.Protocol):
def __init__(self):
self.transport = None

def connection_made(self, transport):
self.transport = transport

def data_received(self, data):
try:
self = weakref.ref(self)
nonlocal is_new
if is_new:
done.set_result(data)
else:
is_new = True
new_proto = Protocol()
self().transport.set_protocol(new_proto)
new_proto.connection_made(self().transport)
new_proto.data_received(data)
except Exception as e:
done.set_exception(e)

async def test(addr):
await self.loop.create_connection(Protocol, *addr)
data = await done
self.assertEqual(data, b'hello')

with self.tcp_server(server) as srv:
self.loop.run_until_complete(test(srv.addr))


class Test_UV_TCP(_TestTCP, tb.UVTestCase):

Expand Down
4 changes: 4 additions & 0 deletions uvloop/loop.pyi
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import asyncio
import ssl
import sys
from socket import AddressFamily, SocketKind, _Address, _RetAddress, socket
from typing import (
IO,
Expand DownExpand Up@@ -210,6 +211,9 @@ class Loop:
async def sock_sendall(self, sock: socket, data: bytes) -> None: ...
async def sock_accept(self, sock: socket) -> Tuple[socket, _RetAddress]: ...
async def sock_connect(self, sock: socket, address: _Address) -> None: ...
async def sock_recvfrom(self, sock: socket, bufsize: int) -> bytes: ...
async def sock_recvfrom_into(self, sock: socket, buf: bytearray, nbytes: int = ...) -> int: ...
async def sock_sendto(self, sock: socket, data: bytes, address: _Address) -> None: ...
async def connect_accepted_socket(
self,
protocol_factory: Callable[[], _ProtocolT],
Expand Down
37 changes: 33 additions & 4 deletions uvloop/loop.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@ include "errors.pyx"

cdef:
int PY39 = PY_VERSION_HEX >= 0x03090000
int PY311 = PY_VERSION_HEX >= 0x030b0000
uint64_t MAX_SLEEP = 3600 * 24 * 365 * 100


Expand DownExpand Up@@ -1413,19 +1414,35 @@ cdef class Loop:
"""Create a Future object attached to the loop."""
return self._new_future()

def create_task(self, coro, *, name=None):
def create_task(self, coro, *, name=None, context=None):
"""Schedule a coroutine object.

Return a task object.

If name is not None, task.set_name(name) will be called if the task
object has the set_name attribute, true for default Task in Python 3.8.

An optional keyword-only context argument allows specifying a custom
contextvars.Context for the coro to run in. The current context copy is
created when no context is provided.
"""
self._check_closed()
if self._task_factory is None:
task = aio_Task(coro, loop=self)
if PY311:
if self._task_factory is None:
task = aio_Task(coro, loop=self, context=context)
else:
task = self._task_factory(self, coro, context=context)
else:
task = self._task_factory(self, coro)
if context is None:
if self._task_factory is None:
task = aio_Task(coro, loop=self)
else:
task = self._task_factory(self, coro)
else:
if self._task_factory is None:
task = context.run(aio_Task, coro, self)
else:
task = context.run(self._task_factory, self, coro)

# copied from asyncio.tasks._set_task_name (bpo-34270)
if name is not None:
Expand DownExpand Up@@ -2604,6 +2621,18 @@ cdef class Loop:
finally:
socket_dec_io_ref(sock)

@cython.iterable_coroutine
async def sock_recvfrom(self, sock, bufsize):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_recvfrom_into(self, sock, buf, nbytes=0):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_sendto(self, sock, data, address):
raise NotImplementedError

@cython.iterable_coroutine
async def connect_accepted_socket(self, protocol_factory, sock, *,
ssl=None,
Expand Down
4 changes: 2 additions & 2 deletions uvloop/pseudosock.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,8 +41,8 @@ cdef class PseudoSocket:

def __repr__(self):
s = ("<uvloop.PseudoSocket fd={}, family={!s}, "
"type={!s}, proto={}").format(self.fileno(), self.family,
self.type, self.proto)
"type={!s}, proto={}").format(self.fileno(), self.family.name,
self.type.name, self.proto)

if self._fd != -1:
try:
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -75,7 +75,7 @@ jobs:
strategy:
matrix:
os: [ubuntu-latest, macos-latest]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*"]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*", "cp311-*"]
cibw_arch: ["x86_64", "aarch64", "universal2"]
exclude:
- os: ubuntu-latest
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/tests.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -14,7 +14,7 @@ jobs:
runs-on: ${{ matrix.os }}
strategy:
matrix:
python-version: ["3.7", "3.8", "3.9", "3.10"]
python-version: ["3.7", "3.8", "3.9", "3.10", "3.11.0-rc.1"]
os: [ubuntu-latest, macos-latest]

env:
Expand Down
6 changes: 4 additions & 2 deletions setup.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,14 +21,16 @@
from setuptools.command.sdist import sdist


CYTHON_DEPENDENCY = 'Cython(>=0.29.24,<0.30.0)'
CYTHON_DEPENDENCY = 'Cython(>=0.29.32,<0.30.0)'

# Minimal dependencies required to test uvloop.
TEST_DEPENDENCIES = [
# pycodestyle is a dependency of flake8, but it must be frozen because
# their combination breaks too often
# (example breakage: https://gitlab.com/pycqa/flake8/issues/427)
'aiohttp',
# aiohttp doesn't support 3.11 yet,
# see https://github.com/aio-libs/aiohttp/issues/6600
'aiohttp ; python_version < "3.11"',
'flake8~=3.9.2',
'psutil',
'pycodestyle~=2.7.0',
Expand Down
8 changes: 6 additions & 2 deletions tests/test_base.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -540,7 +540,9 @@ class MyTask(asyncio.Task):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
self.loop.set_task_factory(factory)
Expand DownExpand Up@@ -577,7 +579,9 @@ def get_name(self):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
task = self.loop.create_task(coro(), name="mytask")
Expand Down
4 changes: 4 additions & 0 deletions tests/test_context.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -452,6 +452,10 @@ def close():
self._run_server_test(test, async_sock=True)

def test_create_ssl_server_manual_connection_lost(self):
if self.implementation == 'asyncio' and sys.version_info >= (3, 11, 0):
# TODO(fantix): fix for 3.11
raise unittest.SkipTest('should pass on 3.11')

async def test(proto, cvar, ssl_sock, **_):
def close():
cvar.set('closing')
Expand Down
37 changes: 0 additions & 37 deletions tests/test_tcp.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -652,43 +652,6 @@ async def runner():
self.assertIsNone(
self.loop.run_until_complete(connection_lost_called))

def test_context_run_segfault(self):
is_new = False
done = self.loop.create_future()

def server(sock):
sock.sendall(b'hello')

class Protocol(asyncio.Protocol):
def __init__(self):
self.transport = None

def connection_made(self, transport):
self.transport = transport

def data_received(self, data):
try:
self = weakref.ref(self)
nonlocal is_new
if is_new:
done.set_result(data)
else:
is_new = True
new_proto = Protocol()
self().transport.set_protocol(new_proto)
new_proto.connection_made(self().transport)
new_proto.data_received(data)
except Exception as e:
done.set_exception(e)

async def test(addr):
await self.loop.create_connection(Protocol, *addr)
data = await done
self.assertEqual(data, b'hello')

with self.tcp_server(server) as srv:
self.loop.run_until_complete(test(srv.addr))


class Test_UV_TCP(_TestTCP, tb.UVTestCase):

Expand Down
4 changes: 4 additions & 0 deletions uvloop/loop.pyi
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import asyncio
import ssl
import sys
from socket import AddressFamily, SocketKind, _Address, _RetAddress, socket
from typing import (
IO,
Expand DownExpand Up@@ -210,6 +211,9 @@ class Loop:
async def sock_sendall(self, sock: socket, data: bytes) -> None: ...
async def sock_accept(self, sock: socket) -> Tuple[socket, _RetAddress]: ...
async def sock_connect(self, sock: socket, address: _Address) -> None: ...
async def sock_recvfrom(self, sock: socket, bufsize: int) -> bytes: ...
async def sock_recvfrom_into(self, sock: socket, buf: bytearray, nbytes: int = ...) -> int: ...
async def sock_sendto(self, sock: socket, data: bytes, address: _Address) -> None: ...
async def connect_accepted_socket(
self,
protocol_factory: Callable[[], _ProtocolT],
Expand Down
37 changes: 33 additions & 4 deletions uvloop/loop.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@ include "errors.pyx"

cdef:
int PY39 = PY_VERSION_HEX >= 0x03090000
int PY311 = PY_VERSION_HEX >= 0x030b0000
uint64_t MAX_SLEEP = 3600 * 24 * 365 * 100


Expand DownExpand Up@@ -1413,19 +1414,35 @@ cdef class Loop:
"""Create a Future object attached to the loop."""
return self._new_future()

def create_task(self, coro, *, name=None):
def create_task(self, coro, *, name=None, context=None):
"""Schedule a coroutine object.

Return a task object.

If name is not None, task.set_name(name) will be called if the task
object has the set_name attribute, true for default Task in Python 3.8.

An optional keyword-only context argument allows specifying a custom
contextvars.Context for the coro to run in. The current context copy is
created when no context is provided.
"""
self._check_closed()
if self._task_factory is None:
task = aio_Task(coro, loop=self)
if PY311:
if self._task_factory is None:
task = aio_Task(coro, loop=self, context=context)
else:
task = self._task_factory(self, coro, context=context)
else:
task = self._task_factory(self, coro)
if context is None:
if self._task_factory is None:
task = aio_Task(coro, loop=self)
else:
task = self._task_factory(self, coro)
else:
if self._task_factory is None:
task = context.run(aio_Task, coro, self)
else:
task = context.run(self._task_factory, self, coro)

# copied from asyncio.tasks._set_task_name (bpo-34270)
if name is not None:
Expand DownExpand Up@@ -2604,6 +2621,18 @@ cdef class Loop:
finally:
socket_dec_io_ref(sock)

@cython.iterable_coroutine
async def sock_recvfrom(self, sock, bufsize):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_recvfrom_into(self, sock, buf, nbytes=0):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_sendto(self, sock, data, address):
raise NotImplementedError

@cython.iterable_coroutine
async def connect_accepted_socket(self, protocol_factory, sock, *,
ssl=None,
Expand Down
4 changes: 2 additions & 2 deletions uvloop/pseudosock.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,8 +41,8 @@ cdef class PseudoSocket:

def __repr__(self):
s = ("<uvloop.PseudoSocket fd={}, family={!s}, "
"type={!s}, proto={}").format(self.fileno(), self.family,
self.type, self.proto)
"type={!s}, proto={}").format(self.fileno(), self.family.name,
self.type.name, self.proto)

if self._fd != -1:
try:
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length \u003e 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -75,7 +75,7 @@ jobs:
strategy:
matrix:
os: [ubuntu-latest, macos-latest]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*"]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*", "cp311-*"]
cibw_arch: ["x86_64", "aarch64", "universal2"]
exclude:
- os: ubuntu-latest
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/tests.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -14,7 +14,7 @@ jobs:
runs-on: ${{ matrix.os }}
strategy:
matrix:
python-version: ["3.7", "3.8", "3.9", "3.10"]
python-version: ["3.7", "3.8", "3.9", "3.10", "3.11.0-rc.1"]
os: [ubuntu-latest, macos-latest]

env:
Expand Down
6 changes: 4 additions & 2 deletions setup.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,14 +21,16 @@
from setuptools.command.sdist import sdist


CYTHON_DEPENDENCY = 'Cython(>=0.29.24,<0.30.0)'
CYTHON_DEPENDENCY = 'Cython(>=0.29.32,<0.30.0)'

# Minimal dependencies required to test uvloop.
TEST_DEPENDENCIES = [
# pycodestyle is a dependency of flake8, but it must be frozen because
# their combination breaks too often
# (example breakage: https://gitlab.com/pycqa/flake8/issues/427)
'aiohttp',
# aiohttp doesn't support 3.11 yet,
# see https://github.com/aio-libs/aiohttp/issues/6600
'aiohttp ; python_version < "3.11"',
'flake8~=3.9.2',
'psutil',
'pycodestyle~=2.7.0',
Expand Down
8 changes: 6 additions & 2 deletions tests/test_base.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -540,7 +540,9 @@ class MyTask(asyncio.Task):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
self.loop.set_task_factory(factory)
Expand DownExpand Up@@ -577,7 +579,9 @@ def get_name(self):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
task = self.loop.create_task(coro(), name="mytask")
Expand Down
4 changes: 4 additions & 0 deletions tests/test_context.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -452,6 +452,10 @@ def close():
self._run_server_test(test, async_sock=True)

def test_create_ssl_server_manual_connection_lost(self):
if self.implementation == 'asyncio' and sys.version_info >= (3, 11, 0):
# TODO(fantix): fix for 3.11
raise unittest.SkipTest('should pass on 3.11')

async def test(proto, cvar, ssl_sock, **_):
def close():
cvar.set('closing')
Expand Down
37 changes: 0 additions & 37 deletions tests/test_tcp.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -652,43 +652,6 @@ async def runner():
self.assertIsNone(
self.loop.run_until_complete(connection_lost_called))

def test_context_run_segfault(self):
is_new = False
done = self.loop.create_future()

def server(sock):
sock.sendall(b'hello')

class Protocol(asyncio.Protocol):
def __init__(self):
self.transport = None

def connection_made(self, transport):
self.transport = transport

def data_received(self, data):
try:
self = weakref.ref(self)
nonlocal is_new
if is_new:
done.set_result(data)
else:
is_new = True
new_proto = Protocol()
self().transport.set_protocol(new_proto)
new_proto.connection_made(self().transport)
new_proto.data_received(data)
except Exception as e:
done.set_exception(e)

async def test(addr):
await self.loop.create_connection(Protocol, *addr)
data = await done
self.assertEqual(data, b'hello')

with self.tcp_server(server) as srv:
self.loop.run_until_complete(test(srv.addr))


class Test_UV_TCP(_TestTCP, tb.UVTestCase):

Expand Down
4 changes: 4 additions & 0 deletions uvloop/loop.pyi
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import asyncio
import ssl
import sys
from socket import AddressFamily, SocketKind, _Address, _RetAddress, socket
from typing import (
IO,
Expand DownExpand Up@@ -210,6 +211,9 @@ class Loop:
async def sock_sendall(self, sock: socket, data: bytes) -> None: ...
async def sock_accept(self, sock: socket) -> Tuple[socket, _RetAddress]: ...
async def sock_connect(self, sock: socket, address: _Address) -> None: ...
async def sock_recvfrom(self, sock: socket, bufsize: int) -> bytes: ...
async def sock_recvfrom_into(self, sock: socket, buf: bytearray, nbytes: int = ...) -> int: ...
async def sock_sendto(self, sock: socket, data: bytes, address: _Address) -> None: ...
async def connect_accepted_socket(
self,
protocol_factory: Callable[[], _ProtocolT],
Expand Down
37 changes: 33 additions & 4 deletions uvloop/loop.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@ include "errors.pyx"

cdef:
int PY39 = PY_VERSION_HEX >= 0x03090000
int PY311 = PY_VERSION_HEX >= 0x030b0000
uint64_t MAX_SLEEP = 3600 * 24 * 365 * 100


Expand DownExpand Up@@ -1413,19 +1414,35 @@ cdef class Loop:
"""Create a Future object attached to the loop."""
return self._new_future()

def create_task(self, coro, *, name=None):
def create_task(self, coro, *, name=None, context=None):
"""Schedule a coroutine object.

Return a task object.

If name is not None, task.set_name(name) will be called if the task
object has the set_name attribute, true for default Task in Python 3.8.

An optional keyword-only context argument allows specifying a custom
contextvars.Context for the coro to run in. The current context copy is
created when no context is provided.
"""
self._check_closed()
if self._task_factory is None:
task = aio_Task(coro, loop=self)
if PY311:
if self._task_factory is None:
task = aio_Task(coro, loop=self, context=context)
else:
task = self._task_factory(self, coro, context=context)
else:
task = self._task_factory(self, coro)
if context is None:
if self._task_factory is None:
task = aio_Task(coro, loop=self)
else:
task = self._task_factory(self, coro)
else:
if self._task_factory is None:
task = context.run(aio_Task, coro, self)
else:
task = context.run(self._task_factory, self, coro)

# copied from asyncio.tasks._set_task_name (bpo-34270)
if name is not None:
Expand DownExpand Up@@ -2604,6 +2621,18 @@ cdef class Loop:
finally:
socket_dec_io_ref(sock)

@cython.iterable_coroutine
async def sock_recvfrom(self, sock, bufsize):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_recvfrom_into(self, sock, buf, nbytes=0):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_sendto(self, sock, data, address):
raise NotImplementedError

@cython.iterable_coroutine
async def connect_accepted_socket(self, protocol_factory, sock, *,
ssl=None,
Expand Down
4 changes: 2 additions & 2 deletions uvloop/pseudosock.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,8 +41,8 @@ cdef class PseudoSocket:

def __repr__(self):
s = ("<uvloop.PseudoSocket fd={}, family={!s}, "
"type={!s}, proto={}").format(self.fileno(), self.family,
self.type, self.proto)
"type={!s}, proto={}").format(self.fileno(), self.family.name,
self.type.name, self.proto)

if self._fd != -1:
try:
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -75,7 +75,7 @@ jobs:
strategy:
matrix:
os: [ubuntu-latest, macos-latest]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*"]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*", "cp311-*"]
cibw_arch: ["x86_64", "aarch64", "universal2"]
exclude:
- os: ubuntu-latest
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/tests.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -14,7 +14,7 @@ jobs:
runs-on: ${{ matrix.os }}
strategy:
matrix:
python-version: ["3.7", "3.8", "3.9", "3.10"]
python-version: ["3.7", "3.8", "3.9", "3.10", "3.11.0-rc.1"]
os: [ubuntu-latest, macos-latest]

env:
Expand Down
6 changes: 4 additions & 2 deletions setup.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,14 +21,16 @@
from setuptools.command.sdist import sdist


CYTHON_DEPENDENCY = 'Cython(>=0.29.24,<0.30.0)'
CYTHON_DEPENDENCY = 'Cython(>=0.29.32,<0.30.0)'

# Minimal dependencies required to test uvloop.
TEST_DEPENDENCIES = [
# pycodestyle is a dependency of flake8, but it must be frozen because
# their combination breaks too often
# (example breakage: https://gitlab.com/pycqa/flake8/issues/427)
'aiohttp',
# aiohttp doesn't support 3.11 yet,
# see https://github.com/aio-libs/aiohttp/issues/6600
'aiohttp ; python_version < "3.11"',
'flake8~=3.9.2',
'psutil',
'pycodestyle~=2.7.0',
Expand Down
8 changes: 6 additions & 2 deletions tests/test_base.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -540,7 +540,9 @@ class MyTask(asyncio.Task):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
self.loop.set_task_factory(factory)
Expand DownExpand Up@@ -577,7 +579,9 @@ def get_name(self):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
task = self.loop.create_task(coro(), name="mytask")
Expand Down
4 changes: 4 additions & 0 deletions tests/test_context.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -452,6 +452,10 @@ def close():
self._run_server_test(test, async_sock=True)

def test_create_ssl_server_manual_connection_lost(self):
if self.implementation == 'asyncio' and sys.version_info >= (3, 11, 0):
# TODO(fantix): fix for 3.11
raise unittest.SkipTest('should pass on 3.11')

async def test(proto, cvar, ssl_sock, **_):
def close():
cvar.set('closing')
Expand Down
37 changes: 0 additions & 37 deletions tests/test_tcp.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -652,43 +652,6 @@ async def runner():
self.assertIsNone(
self.loop.run_until_complete(connection_lost_called))

def test_context_run_segfault(self):
is_new = False
done = self.loop.create_future()

def server(sock):
sock.sendall(b'hello')

class Protocol(asyncio.Protocol):
def __init__(self):
self.transport = None

def connection_made(self, transport):
self.transport = transport

def data_received(self, data):
try:
self = weakref.ref(self)
nonlocal is_new
if is_new:
done.set_result(data)
else:
is_new = True
new_proto = Protocol()
self().transport.set_protocol(new_proto)
new_proto.connection_made(self().transport)
new_proto.data_received(data)
except Exception as e:
done.set_exception(e)

async def test(addr):
await self.loop.create_connection(Protocol, *addr)
data = await done
self.assertEqual(data, b'hello')

with self.tcp_server(server) as srv:
self.loop.run_until_complete(test(srv.addr))


class Test_UV_TCP(_TestTCP, tb.UVTestCase):

Expand Down
4 changes: 4 additions & 0 deletions uvloop/loop.pyi
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import asyncio
import ssl
import sys
from socket import AddressFamily, SocketKind, _Address, _RetAddress, socket
from typing import (
IO,
Expand DownExpand Up@@ -210,6 +211,9 @@ class Loop:
async def sock_sendall(self, sock: socket, data: bytes) -> None: ...
async def sock_accept(self, sock: socket) -> Tuple[socket, _RetAddress]: ...
async def sock_connect(self, sock: socket, address: _Address) -> None: ...
async def sock_recvfrom(self, sock: socket, bufsize: int) -> bytes: ...
async def sock_recvfrom_into(self, sock: socket, buf: bytearray, nbytes: int = ...) -> int: ...
async def sock_sendto(self, sock: socket, data: bytes, address: _Address) -> None: ...
async def connect_accepted_socket(
self,
protocol_factory: Callable[[], _ProtocolT],
Expand Down
37 changes: 33 additions & 4 deletions uvloop/loop.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@ include "errors.pyx"

cdef:
int PY39 = PY_VERSION_HEX >= 0x03090000
int PY311 = PY_VERSION_HEX >= 0x030b0000
uint64_t MAX_SLEEP = 3600 * 24 * 365 * 100


Expand DownExpand Up@@ -1413,19 +1414,35 @@ cdef class Loop:
"""Create a Future object attached to the loop."""
return self._new_future()

def create_task(self, coro, *, name=None):
def create_task(self, coro, *, name=None, context=None):
"""Schedule a coroutine object.

Return a task object.

If name is not None, task.set_name(name) will be called if the task
object has the set_name attribute, true for default Task in Python 3.8.

An optional keyword-only context argument allows specifying a custom
contextvars.Context for the coro to run in. The current context copy is
created when no context is provided.
"""
self._check_closed()
if self._task_factory is None:
task = aio_Task(coro, loop=self)
if PY311:
if self._task_factory is None:
task = aio_Task(coro, loop=self, context=context)
else:
task = self._task_factory(self, coro, context=context)
else:
task = self._task_factory(self, coro)
if context is None:
if self._task_factory is None:
task = aio_Task(coro, loop=self)
else:
task = self._task_factory(self, coro)
else:
if self._task_factory is None:
task = context.run(aio_Task, coro, self)
else:
task = context.run(self._task_factory, self, coro)

# copied from asyncio.tasks._set_task_name (bpo-34270)
if name is not None:
Expand DownExpand Up@@ -2604,6 +2621,18 @@ cdef class Loop:
finally:
socket_dec_io_ref(sock)

@cython.iterable_coroutine
async def sock_recvfrom(self, sock, bufsize):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_recvfrom_into(self, sock, buf, nbytes=0):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_sendto(self, sock, data, address):
raise NotImplementedError

@cython.iterable_coroutine
async def connect_accepted_socket(self, protocol_factory, sock, *,
ssl=None,
Expand Down
4 changes: 2 additions & 2 deletions uvloop/pseudosock.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,8 +41,8 @@ cdef class PseudoSocket:

def __repr__(self):
s = ("<uvloop.PseudoSocket fd={}, family={!s}, "
"type={!s}, proto={}").format(self.fileno(), self.family,
self.type, self.proto)
"type={!s}, proto={}").format(self.fileno(), self.family.name,
self.type.name, self.proto)

if self._fd != -1:
try:
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -75,7 +75,7 @@ jobs:
strategy:
matrix:
os: [ubuntu-latest, macos-latest]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*"]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*", "cp311-*"]
cibw_arch: ["x86_64", "aarch64", "universal2"]
exclude:
- os: ubuntu-latest
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/tests.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -14,7 +14,7 @@ jobs:
runs-on: ${{ matrix.os }}
strategy:
matrix:
python-version: ["3.7", "3.8", "3.9", "3.10"]
python-version: ["3.7", "3.8", "3.9", "3.10", "3.11.0-rc.1"]
os: [ubuntu-latest, macos-latest]

env:
Expand Down
6 changes: 4 additions & 2 deletions setup.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,14 +21,16 @@
from setuptools.command.sdist import sdist


CYTHON_DEPENDENCY = 'Cython(>=0.29.24,<0.30.0)'
CYTHON_DEPENDENCY = 'Cython(>=0.29.32,<0.30.0)'

# Minimal dependencies required to test uvloop.
TEST_DEPENDENCIES = [
# pycodestyle is a dependency of flake8, but it must be frozen because
# their combination breaks too often
# (example breakage: https://gitlab.com/pycqa/flake8/issues/427)
'aiohttp',
# aiohttp doesn't support 3.11 yet,
# see https://github.com/aio-libs/aiohttp/issues/6600
'aiohttp ; python_version < "3.11"',
'flake8~=3.9.2',
'psutil',
'pycodestyle~=2.7.0',
Expand Down
8 changes: 6 additions & 2 deletions tests/test_base.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -540,7 +540,9 @@ class MyTask(asyncio.Task):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
self.loop.set_task_factory(factory)
Expand DownExpand Up@@ -577,7 +579,9 @@ def get_name(self):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
task = self.loop.create_task(coro(), name="mytask")
Expand Down
4 changes: 4 additions & 0 deletions tests/test_context.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -452,6 +452,10 @@ def close():
self._run_server_test(test, async_sock=True)

def test_create_ssl_server_manual_connection_lost(self):
if self.implementation == 'asyncio' and sys.version_info >= (3, 11, 0):
# TODO(fantix): fix for 3.11
raise unittest.SkipTest('should pass on 3.11')

async def test(proto, cvar, ssl_sock, **_):
def close():
cvar.set('closing')
Expand Down
37 changes: 0 additions & 37 deletions tests/test_tcp.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -652,43 +652,6 @@ async def runner():
self.assertIsNone(
self.loop.run_until_complete(connection_lost_called))

def test_context_run_segfault(self):
is_new = False
done = self.loop.create_future()

def server(sock):
sock.sendall(b'hello')

class Protocol(asyncio.Protocol):
def __init__(self):
self.transport = None

def connection_made(self, transport):
self.transport = transport

def data_received(self, data):
try:
self = weakref.ref(self)
nonlocal is_new
if is_new:
done.set_result(data)
else:
is_new = True
new_proto = Protocol()
self().transport.set_protocol(new_proto)
new_proto.connection_made(self().transport)
new_proto.data_received(data)
except Exception as e:
done.set_exception(e)

async def test(addr):
await self.loop.create_connection(Protocol, *addr)
data = await done
self.assertEqual(data, b'hello')

with self.tcp_server(server) as srv:
self.loop.run_until_complete(test(srv.addr))


class Test_UV_TCP(_TestTCP, tb.UVTestCase):

Expand Down
4 changes: 4 additions & 0 deletions uvloop/loop.pyi
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import asyncio
import ssl
import sys
from socket import AddressFamily, SocketKind, _Address, _RetAddress, socket
from typing import (
IO,
Expand DownExpand Up@@ -210,6 +211,9 @@ class Loop:
async def sock_sendall(self, sock: socket, data: bytes) -> None: ...
async def sock_accept(self, sock: socket) -> Tuple[socket, _RetAddress]: ...
async def sock_connect(self, sock: socket, address: _Address) -> None: ...
async def sock_recvfrom(self, sock: socket, bufsize: int) -> bytes: ...
async def sock_recvfrom_into(self, sock: socket, buf: bytearray, nbytes: int = ...) -> int: ...
async def sock_sendto(self, sock: socket, data: bytes, address: _Address) -> None: ...
async def connect_accepted_socket(
self,
protocol_factory: Callable[[], _ProtocolT],
Expand Down
37 changes: 33 additions & 4 deletions uvloop/loop.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@ include "errors.pyx"

cdef:
int PY39 = PY_VERSION_HEX >= 0x03090000
int PY311 = PY_VERSION_HEX >= 0x030b0000
uint64_t MAX_SLEEP = 3600 * 24 * 365 * 100


Expand DownExpand Up@@ -1413,19 +1414,35 @@ cdef class Loop:
"""Create a Future object attached to the loop."""
return self._new_future()

def create_task(self, coro, *, name=None):
def create_task(self, coro, *, name=None, context=None):
"""Schedule a coroutine object.

Return a task object.

If name is not None, task.set_name(name) will be called if the task
object has the set_name attribute, true for default Task in Python 3.8.

An optional keyword-only context argument allows specifying a custom
contextvars.Context for the coro to run in. The current context copy is
created when no context is provided.
"""
self._check_closed()
if self._task_factory is None:
task = aio_Task(coro, loop=self)
if PY311:
if self._task_factory is None:
task = aio_Task(coro, loop=self, context=context)
else:
task = self._task_factory(self, coro, context=context)
else:
task = self._task_factory(self, coro)
if context is None:
if self._task_factory is None:
task = aio_Task(coro, loop=self)
else:
task = self._task_factory(self, coro)
else:
if self._task_factory is None:
task = context.run(aio_Task, coro, self)
else:
task = context.run(self._task_factory, self, coro)

# copied from asyncio.tasks._set_task_name (bpo-34270)
if name is not None:
Expand DownExpand Up@@ -2604,6 +2621,18 @@ cdef class Loop:
finally:
socket_dec_io_ref(sock)

@cython.iterable_coroutine
async def sock_recvfrom(self, sock, bufsize):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_recvfrom_into(self, sock, buf, nbytes=0):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_sendto(self, sock, data, address):
raise NotImplementedError

@cython.iterable_coroutine
async def connect_accepted_socket(self, protocol_factory, sock, *,
ssl=None,
Expand Down
4 changes: 2 additions & 2 deletions uvloop/pseudosock.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,8 +41,8 @@ cdef class PseudoSocket:

def __repr__(self):
s = ("<uvloop.PseudoSocket fd={}, family={!s}, "
"type={!s}, proto={}").format(self.fileno(), self.family,
self.type, self.proto)
"type={!s}, proto={}").format(self.fileno(), self.family.name,
self.type.name, self.proto)

if self._fd != -1:
try:
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -75,7 +75,7 @@ jobs:
strategy:
matrix:
os: [ubuntu-latest, macos-latest]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*"]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*", "cp311-*"]
cibw_arch: ["x86_64", "aarch64", "universal2"]
exclude:
- os: ubuntu-latest
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/tests.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -14,7 +14,7 @@ jobs:
runs-on: ${{ matrix.os }}
strategy:
matrix:
python-version: ["3.7", "3.8", "3.9", "3.10"]
python-version: ["3.7", "3.8", "3.9", "3.10", "3.11.0-rc.1"]
os: [ubuntu-latest, macos-latest]

env:
Expand Down
6 changes: 4 additions & 2 deletions setup.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,14 +21,16 @@
from setuptools.command.sdist import sdist


CYTHON_DEPENDENCY = 'Cython(>=0.29.24,<0.30.0)'
CYTHON_DEPENDENCY = 'Cython(>=0.29.32,<0.30.0)'

# Minimal dependencies required to test uvloop.
TEST_DEPENDENCIES = [
# pycodestyle is a dependency of flake8, but it must be frozen because
# their combination breaks too often
# (example breakage: https://gitlab.com/pycqa/flake8/issues/427)
'aiohttp',
# aiohttp doesn't support 3.11 yet,
# see https://github.com/aio-libs/aiohttp/issues/6600
'aiohttp ; python_version < "3.11"',
'flake8~=3.9.2',
'psutil',
'pycodestyle~=2.7.0',
Expand Down
8 changes: 6 additions & 2 deletions tests/test_base.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -540,7 +540,9 @@ class MyTask(asyncio.Task):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
self.loop.set_task_factory(factory)
Expand DownExpand Up@@ -577,7 +579,9 @@ def get_name(self):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
task = self.loop.create_task(coro(), name="mytask")
Expand Down
4 changes: 4 additions & 0 deletions tests/test_context.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -452,6 +452,10 @@ def close():
self._run_server_test(test, async_sock=True)

def test_create_ssl_server_manual_connection_lost(self):
if self.implementation == 'asyncio' and sys.version_info >= (3, 11, 0):
# TODO(fantix): fix for 3.11
raise unittest.SkipTest('should pass on 3.11')

async def test(proto, cvar, ssl_sock, **_):
def close():
cvar.set('closing')
Expand Down
37 changes: 0 additions & 37 deletions tests/test_tcp.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -652,43 +652,6 @@ async def runner():
self.assertIsNone(
self.loop.run_until_complete(connection_lost_called))

def test_context_run_segfault(self):
is_new = False
done = self.loop.create_future()

def server(sock):
sock.sendall(b'hello')

class Protocol(asyncio.Protocol):
def __init__(self):
self.transport = None

def connection_made(self, transport):
self.transport = transport

def data_received(self, data):
try:
self = weakref.ref(self)
nonlocal is_new
if is_new:
done.set_result(data)
else:
is_new = True
new_proto = Protocol()
self().transport.set_protocol(new_proto)
new_proto.connection_made(self().transport)
new_proto.data_received(data)
except Exception as e:
done.set_exception(e)

async def test(addr):
await self.loop.create_connection(Protocol, *addr)
data = await done
self.assertEqual(data, b'hello')

with self.tcp_server(server) as srv:
self.loop.run_until_complete(test(srv.addr))


class Test_UV_TCP(_TestTCP, tb.UVTestCase):

Expand Down
4 changes: 4 additions & 0 deletions uvloop/loop.pyi
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import asyncio
import ssl
import sys
from socket import AddressFamily, SocketKind, _Address, _RetAddress, socket
from typing import (
IO,
Expand DownExpand Up@@ -210,6 +211,9 @@ class Loop:
async def sock_sendall(self, sock: socket, data: bytes) -> None: ...
async def sock_accept(self, sock: socket) -> Tuple[socket, _RetAddress]: ...
async def sock_connect(self, sock: socket, address: _Address) -> None: ...
async def sock_recvfrom(self, sock: socket, bufsize: int) -> bytes: ...
async def sock_recvfrom_into(self, sock: socket, buf: bytearray, nbytes: int = ...) -> int: ...
async def sock_sendto(self, sock: socket, data: bytes, address: _Address) -> None: ...
async def connect_accepted_socket(
self,
protocol_factory: Callable[[], _ProtocolT],
Expand Down
37 changes: 33 additions & 4 deletions uvloop/loop.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@ include "errors.pyx"

cdef:
int PY39 = PY_VERSION_HEX >= 0x03090000
int PY311 = PY_VERSION_HEX >= 0x030b0000
uint64_t MAX_SLEEP = 3600 * 24 * 365 * 100


Expand DownExpand Up@@ -1413,19 +1414,35 @@ cdef class Loop:
"""Create a Future object attached to the loop."""
return self._new_future()

def create_task(self, coro, *, name=None):
def create_task(self, coro, *, name=None, context=None):
"""Schedule a coroutine object.

Return a task object.

If name is not None, task.set_name(name) will be called if the task
object has the set_name attribute, true for default Task in Python 3.8.

An optional keyword-only context argument allows specifying a custom
contextvars.Context for the coro to run in. The current context copy is
created when no context is provided.
"""
self._check_closed()
if self._task_factory is None:
task = aio_Task(coro, loop=self)
if PY311:
if self._task_factory is None:
task = aio_Task(coro, loop=self, context=context)
else:
task = self._task_factory(self, coro, context=context)
else:
task = self._task_factory(self, coro)
if context is None:
if self._task_factory is None:
task = aio_Task(coro, loop=self)
else:
task = self._task_factory(self, coro)
else:
if self._task_factory is None:
task = context.run(aio_Task, coro, self)
else:
task = context.run(self._task_factory, self, coro)

# copied from asyncio.tasks._set_task_name (bpo-34270)
if name is not None:
Expand DownExpand Up@@ -2604,6 +2621,18 @@ cdef class Loop:
finally:
socket_dec_io_ref(sock)

@cython.iterable_coroutine
async def sock_recvfrom(self, sock, bufsize):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_recvfrom_into(self, sock, buf, nbytes=0):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_sendto(self, sock, data, address):
raise NotImplementedError

@cython.iterable_coroutine
async def connect_accepted_socket(self, protocol_factory, sock, *,
ssl=None,
Expand Down
4 changes: 2 additions & 2 deletions uvloop/pseudosock.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,8 +41,8 @@ cdef class PseudoSocket:

def __repr__(self):
s = ("<uvloop.PseudoSocket fd={}, family={!s}, "
"type={!s}, proto={}").format(self.fileno(), self.family,
self.type, self.proto)
"type={!s}, proto={}").format(self.fileno(), self.family.name,
self.type.name, self.proto)

if self._fd != -1:
try:
Expand Down
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -75,7 +75,7 @@ jobs:
strategy:
matrix:
os: [ubuntu-latest, macos-latest]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*"]
cibw_python: ["cp37-*", "cp38-*", "cp39-*", "cp310-*", "cp311-*"]
cibw_arch: ["x86_64", "aarch64", "universal2"]
exclude:
- os: ubuntu-latest
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/tests.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -14,7 +14,7 @@ jobs:
runs-on: ${{ matrix.os }}
strategy:
matrix:
python-version: ["3.7", "3.8", "3.9", "3.10"]
python-version: ["3.7", "3.8", "3.9", "3.10", "3.11.0-rc.1"]
os: [ubuntu-latest, macos-latest]

env:
Expand Down
6 changes: 4 additions & 2 deletions setup.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,14 +21,16 @@
from setuptools.command.sdist import sdist


CYTHON_DEPENDENCY = 'Cython(>=0.29.24,<0.30.0)'
CYTHON_DEPENDENCY = 'Cython(>=0.29.32,<0.30.0)'

# Minimal dependencies required to test uvloop.
TEST_DEPENDENCIES = [
# pycodestyle is a dependency of flake8, but it must be frozen because
# their combination breaks too often
# (example breakage: https://gitlab.com/pycqa/flake8/issues/427)
'aiohttp',
# aiohttp doesn't support 3.11 yet,
# see https://github.com/aio-libs/aiohttp/issues/6600
'aiohttp ; python_version < "3.11"',
'flake8~=3.9.2',
'psutil',
'pycodestyle~=2.7.0',
Expand Down
8 changes: 6 additions & 2 deletions tests/test_base.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -540,7 +540,9 @@ class MyTask(asyncio.Task):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
self.loop.set_task_factory(factory)
Expand DownExpand Up@@ -577,7 +579,9 @@ def get_name(self):
async def coro():
pass

factory = lambda loop, coro: MyTask(coro, loop=loop)
factory = lambda loop, coro, **kwargs: MyTask(
coro, loop=loop, **kwargs
)

self.assertIsNone(self.loop.get_task_factory())
task = self.loop.create_task(coro(), name="mytask")
Expand Down
4 changes: 4 additions & 0 deletions tests/test_context.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -452,6 +452,10 @@ def close():
self._run_server_test(test, async_sock=True)

def test_create_ssl_server_manual_connection_lost(self):
if self.implementation == 'asyncio' and sys.version_info >= (3, 11, 0):
# TODO(fantix): fix for 3.11
raise unittest.SkipTest('should pass on 3.11')

async def test(proto, cvar, ssl_sock, **_):
def close():
cvar.set('closing')
Expand Down
37 changes: 0 additions & 37 deletions tests/test_tcp.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -652,43 +652,6 @@ async def runner():
self.assertIsNone(
self.loop.run_until_complete(connection_lost_called))

def test_context_run_segfault(self):
is_new = False
done = self.loop.create_future()

def server(sock):
sock.sendall(b'hello')

class Protocol(asyncio.Protocol):
def __init__(self):
self.transport = None

def connection_made(self, transport):
self.transport = transport

def data_received(self, data):
try:
self = weakref.ref(self)
nonlocal is_new
if is_new:
done.set_result(data)
else:
is_new = True
new_proto = Protocol()
self().transport.set_protocol(new_proto)
new_proto.connection_made(self().transport)
new_proto.data_received(data)
except Exception as e:
done.set_exception(e)

async def test(addr):
await self.loop.create_connection(Protocol, *addr)
data = await done
self.assertEqual(data, b'hello')

with self.tcp_server(server) as srv:
self.loop.run_until_complete(test(srv.addr))


class Test_UV_TCP(_TestTCP, tb.UVTestCase):

Expand Down
4 changes: 4 additions & 0 deletions uvloop/loop.pyi
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
import asyncio
import ssl
import sys
from socket import AddressFamily, SocketKind, _Address, _RetAddress, socket
from typing import (
IO,
Expand DownExpand Up@@ -210,6 +211,9 @@ class Loop:
async def sock_sendall(self, sock: socket, data: bytes) -> None: ...
async def sock_accept(self, sock: socket) -> Tuple[socket, _RetAddress]: ...
async def sock_connect(self, sock: socket, address: _Address) -> None: ...
async def sock_recvfrom(self, sock: socket, bufsize: int) -> bytes: ...
async def sock_recvfrom_into(self, sock: socket, buf: bytearray, nbytes: int = ...) -> int: ...
async def sock_sendto(self, sock: socket, data: bytes, address: _Address) -> None: ...
async def connect_accepted_socket(
self,
protocol_factory: Callable[[], _ProtocolT],
Expand Down
37 changes: 33 additions & 4 deletions uvloop/loop.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -50,6 +50,7 @@ include "errors.pyx"

cdef:
int PY39 = PY_VERSION_HEX >= 0x03090000
int PY311 = PY_VERSION_HEX >= 0x030b0000
uint64_t MAX_SLEEP = 3600 * 24 * 365 * 100


Expand DownExpand Up@@ -1413,19 +1414,35 @@ cdef class Loop:
"""Create a Future object attached to the loop."""
return self._new_future()

def create_task(self, coro, *, name=None):
def create_task(self, coro, *, name=None, context=None):
"""Schedule a coroutine object.

Return a task object.

If name is not None, task.set_name(name) will be called if the task
object has the set_name attribute, true for default Task in Python 3.8.

An optional keyword-only context argument allows specifying a custom
contextvars.Context for the coro to run in. The current context copy is
created when no context is provided.
"""
self._check_closed()
if self._task_factory is None:
task = aio_Task(coro, loop=self)
if PY311:
if self._task_factory is None:
task = aio_Task(coro, loop=self, context=context)
else:
task = self._task_factory(self, coro, context=context)
else:
task = self._task_factory(self, coro)
if context is None:
if self._task_factory is None:
task = aio_Task(coro, loop=self)
else:
task = self._task_factory(self, coro)
else:
if self._task_factory is None:
task = context.run(aio_Task, coro, self)
else:
task = context.run(self._task_factory, self, coro)

# copied from asyncio.tasks._set_task_name (bpo-34270)
if name is not None:
Expand DownExpand Up@@ -2604,6 +2621,18 @@ cdef class Loop:
finally:
socket_dec_io_ref(sock)

@cython.iterable_coroutine
async def sock_recvfrom(self, sock, bufsize):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_recvfrom_into(self, sock, buf, nbytes=0):
raise NotImplementedError

@cython.iterable_coroutine
async def sock_sendto(self, sock, data, address):
raise NotImplementedError

@cython.iterable_coroutine
async def connect_accepted_socket(self, protocol_factory, sock, *,
ssl=None,
Expand Down
4 changes: 2 additions & 2 deletions uvloop/pseudosock.pyx
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,8 +41,8 @@ cdef class PseudoSocket:

def __repr__(self):
s = ("<uvloop.PseudoSocket fd={}, family={!s}, "
"type={!s}, proto={}").format(self.fileno(), self.family,
self.type, self.proto)
"type={!s}, proto={}").format(self.fileno(), self.family.name,
self.type.name, self.proto)

if self._fd != -1:
try:
Expand Down