From 6318f4a40d852b3cba2d7ce755e5e93b79cb1fbe Mon Sep 17 00:00:00 2001 From: David Brochart Date: Fri, 24 Jul 2026 11:40:34 +0200 Subject: [PATCH] Update dependencies --- plugins/cell/pyproject.toml | 2 +- plugins/cell/txl_cell/main.py | 33 +++++++------ plugins/console/pyproject.toml | 2 +- plugins/kernel/pyproject.toml | 2 +- plugins/kernel/txl_kernel/driver.py | 48 ++++++++++++------- plugins/local_contents/pyproject.toml | 2 +- plugins/local_kernels/pyproject.toml | 2 +- .../local_kernels/txl_local_kernels/driver.py | 12 ++--- plugins/local_terminals/pyproject.toml | 2 +- .../txl_local_terminals/terminal.py | 24 ++++++---- .../txl_local_terminals/win_terminal.py | 23 +++++---- .../txl_markdown_viewer/main.py | 3 +- plugins/notebook_editor/pyproject.toml | 2 +- .../txl_notebook_editor/main.py | 21 ++++---- plugins/remote_contents/pyproject.toml | 2 +- plugins/remote_kernels/pyproject.toml | 2 +- .../txl_remote_kernels/driver.py | 7 ++- .../txl_remote_terminals/main.py | 21 ++++---- plugins/terminal/txl_terminal/main.py | 11 ++--- plugins/text_editor/txl_text_editor/main.py | 23 +++++---- plugins/widgets/pyproject.toml | 2 +- txl/pyproject.toml | 2 +- txl/txl/stapled.py | 11 +++++ 23 files changed, 156 insertions(+), 103 deletions(-) create mode 100644 txl/txl/stapled.py diff --git a/plugins/cell/pyproject.toml b/plugins/cell/pyproject.toml index c4ec40f..d600225 100644 --- a/plugins/cell/pyproject.toml +++ b/plugins/cell/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "pycrdt >=0.12.44,<0.13.0",] +dependencies = [ "txl ==0.3.3", "pycrdt >=0.14.1,<0.15.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/cell/txl_cell/main.py b/plugins/cell/txl_cell/main.py index 2830452..6b47997 100644 --- a/plugins/cell/txl_cell/main.py +++ b/plugins/cell/txl_cell/main.py @@ -1,9 +1,10 @@ import json +import math from functools import partial from importlib.metadata import entry_points +from typing import Any -from anyio import create_task_group, sleep -from anyioutils import Queue, create_task +from anyio import create_memory_object_stream, create_task_group, sleep from fps import Module from pycrdt import Doc, Map, MapEvent, Text from rich.text import Text as RichText @@ -12,6 +13,7 @@ from textual.widgets import Static from txl.base import Cell, CellFactory, Contents, Kernel, Widgets +from txl.stapled import StapledObjectStream from txl.text_input import TextInput YDOCS = {ep.name: ep.load() for ep in entry_points(group="ypywidgets")} @@ -75,23 +77,27 @@ def __init__( self.update(mount=False) self.ycell.observe_deep(self.on_change) self.styles.height = "auto" - self.cell_change_events = Queue() - self.widget_change_events = Queue() - create_task(self.observe_cell_changes(), task_group) - create_task(self.observe_widget_changes(), task_group) + self.cell_change_events = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) + self.widget_change_events = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) + task_group.create_task(self.observe_cell_changes()) + task_group.create_task(self.observe_widget_changes()) def on_click(self): self.clicked = True def on_change(self, events): - self.cell_change_events.put_nowait(events) + self.cell_change_events.send_nowait(events) def on_widget_change(self, ydoc, event): - self.widget_change_events.put_nowait((ydoc, event)) + self.widget_change_events.send_nowait((ydoc, event)) async def observe_widget_changes(self): while True: - ydoc, event = await self.widget_change_events.get() + ydoc, event = await self.widget_change_events.receive() model_name = event.delta[0]["insert"] model = YDOCS[f"{model_name}Model"](ydoc=ydoc) widget = YDOCS[f"txl_{model_name}"](model) @@ -100,7 +106,7 @@ async def observe_widget_changes(self): async def observe_cell_changes(self): while True: - events = await self.cell_change_events.get() + events = await self.cell_change_events.receive() for event in events: if isinstance(event, MapEvent): if "execution_state" in event.keys: @@ -150,9 +156,8 @@ async def observe_cell_changes(self): # this is a widget is_widget = True room_id = f"ywidget:{inserted.guid}" - create_task( + self.task_group.create_task( self.contents.websocket_provider(room_id, inserted), - self.task_group, ) inserted["_model_name"] = model_name = Text() model_name.observe(partial( @@ -211,7 +216,7 @@ def update(self, mount: bool = True): language=language, show_border=self.show_border, ) - create_task(self.source.start(), self.task_group) + self.task_group.create_task(self.source.start()) if mount: self.mount(self.source) @@ -227,7 +232,7 @@ def get_output_widget(self, output): guid = output["guid"] ywidget_doc = Doc() room_id = f"ywidget:{guid}" - create_task(self.contents.websocket_provider(room_id, ywidget_doc), self.task_group) + self.task_group.create_task(self.contents.websocket_provider(room_id, ywidget_doc)) ywidget_doc["_model_name"] = model_name = Text() model_name.observe(partial(self.on_widget_change, ywidget_doc)) return diff --git a/plugins/console/pyproject.toml b/plugins/console/pyproject.toml index 6719542..03b0a46 100644 --- a/plugins/console/pyproject.toml +++ b/plugins/console/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "jupyter-ydoc >=3.0.2,<4.0.0",] +dependencies = [ "txl ==0.3.3", "jupyter-ydoc >=4.1.1,<5.0.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/kernel/pyproject.toml b/plugins/kernel/pyproject.toml index 960185e..d585ae5 100644 --- a/plugins/kernel/pyproject.toml +++ b/plugins/kernel/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "python-dateutil >=2.8.2", "pycrdt >=0.12.44,<0.13.0",] +dependencies = [ "txl ==0.3.3", "python-dateutil >=2.8.2", "pycrdt >=0.14.1,<0.15.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/kernel/txl_kernel/driver.py b/plugins/kernel/txl_kernel/driver.py index 2cd4ec9..c006da4 100644 --- a/plugins/kernel/txl_kernel/driver.py +++ b/plugins/kernel/txl_kernel/driver.py @@ -1,11 +1,13 @@ +import math import time -from typing import Dict +from typing import Any, Dict -from anyio import move_on_after -from anyioutils import Event, Queue, create_task +from anyio import Event, create_memory_object_stream, move_on_after from fps import Signal from pycrdt import Array, Map +from txl.stapled import StapledObjectStream + from .message import create_message @@ -31,8 +33,8 @@ def send(self, buffers): buffers=buffers, ) self.msg_cnt += 1 - create_task( - self.send_message(msg, self.shell_channel, change_date_to_str=True), self.task_group + self.task_group.create_task( + self.send_message(msg, self.shell_channel, change_date_to_str=True) ) @@ -41,10 +43,12 @@ def __init__(self, task_group): self.task_group = task_group self.busy = Signal[bool]() self.msg_cnt = 0 - self.execute_requests: Dict[str, Dict[str, Queue]] = {} - self.recv_queue = Queue() + self.execute_requests: Dict[str, Dict[str, StapledObjectStream]] = {} + self.recv_queue = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) self.started = Event() - create_task(self.recv(), task_group) + task_group.create_task(self.recv()) def create_message(self, *args, **kwargs): return create_message(*args, **kwargs) @@ -62,18 +66,22 @@ async def wait_for_ready(self, timeout=float("inf")): await self.send_message(msg, self.shell_channel, change_date_to_str=True) msg_id = msg["header"]["msg_id"] self.execute_requests[msg_id] = { - "iopub": Queue(), - "shell": Queue(), + "iopub": StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ), + "shell": StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ), } with move_on_after(new_timeout) as scope: - msg = await self.execute_requests[msg_id]["shell"].get() + msg = await self.execute_requests[msg_id]["shell"].receive() if scope.cancelled_caught: del self.execute_requests[msg_id] error_message = f"Kernel didn't respond in {timeout} seconds" raise RuntimeError(error_message) if msg["header"]["msg_type"] == "kernel_info_reply": with move_on_after(0.2) as scope: - msg = await self.execute_requests[msg_id]["iopub"].get() + msg = await self.execute_requests[msg_id]["iopub"].receive() if not scope.cancelled_caught: break del self.execute_requests[msg_id] @@ -81,7 +89,7 @@ async def wait_for_ready(self, timeout=float("inf")): async def recv(self): while True: - msg = await self.recv_queue.get() + msg = await self.recv_queue.receive() channel = msg.pop("channel") msg_type = msg["header"]["msg_type"] if msg_type == "comm_open": @@ -108,7 +116,7 @@ async def recv(self): if msg_id in self.execute_requests: # msg["header"] = str_to_date(msg["header"]) # msg["parent_header"] = str_to_date(msg["parent_header"]) - self.execute_requests[msg_id][channel].put_nowait(msg) + self.execute_requests[msg_id][channel].send_nowait(msg) async def execute( self, @@ -136,15 +144,19 @@ async def execute( msg_id = msg["header"]["msg_id"] self.msg_cnt += 1 self.execute_requests[msg_id] = { - "iopub": Queue(), - "shell": Queue(), + "iopub": StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ), + "shell": StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ), } await self.send_message(msg, self.shell_channel, change_date_to_str=True) if wait_for_executed: deadline = time.monotonic() + timeout while True: with move_on_after(deadline_to_timeout(deadline)) as scope: - msg = await self.execute_requests[msg_id]["iopub"].get() + msg = await self.execute_requests[msg_id]["iopub"].receive() if scope.cancelled_caught: del self.execute_requests[msg_id] error_message = f"Kernel didn't respond in {timeout} seconds" @@ -156,7 +168,7 @@ async def execute( ): break with move_on_after(deadline_to_timeout(deadline)) as scope: - msg = await self.execute_requests[msg_id]["shell"].get() + msg = await self.execute_requests[msg_id]["shell"].receive() if scope.cancelled_caught: del self.execute_requests[msg_id] error_message = f"Kernel didn't respond in {timeout} seconds" diff --git a/plugins/local_contents/pyproject.toml b/plugins/local_contents/pyproject.toml index faf6dba..55100d6 100644 --- a/plugins/local_contents/pyproject.toml +++ b/plugins/local_contents/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "anyio >=3.7.0,<5", "jupyter-ydoc >=3.0.2,<4.0.0",] +dependencies = [ "txl ==0.3.3", "anyio >=4.14.2,<5", "jupyter-ydoc >=4.1.1,<5.0.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/local_kernels/pyproject.toml b/plugins/local_kernels/pyproject.toml index bf7dcec..d63abde 100644 --- a/plugins/local_kernels/pyproject.toml +++ b/plugins/local_kernels/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "txl_kernel", "pyzmq >=24.0.1", "ipykernel",] +dependencies = [ "txl ==0.3.3", "txl_kernel", "pyzmq >=27.1.0", "ipykernel",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/local_kernels/txl_local_kernels/driver.py b/plugins/local_kernels/txl_local_kernels/driver.py index d1fa32e..a04ea14 100644 --- a/plugins/local_kernels/txl_local_kernels/driver.py +++ b/plugins/local_kernels/txl_local_kernels/driver.py @@ -2,7 +2,7 @@ import uuid from typing import Any, Dict, List, Optional, cast -from anyioutils import Task, create_task +from anyio import TaskHandle from txl_kernel.driver import KernelMixin from .connect import cfg_t, connect_channel, launch_kernel, read_connection_file @@ -39,7 +39,7 @@ def __init__( self.connection_cfg = read_connection_file(connection_file) self.key = cast(str, self.connection_cfg["key"]) self.session_id = uuid.uuid4().hex - self.channel_tasks: List[Task] = [] + self.channel_tasks: List[TaskHandle] = [] self.comm_handlers = comm_handlers task_group.start_soon(self.start) kernel_drivers.append(self) @@ -83,8 +83,8 @@ def connect_channels(self, connection_cfg: Optional[cfg_t] = None): self.iopub_channel = connect_channel("iopub", connection_cfg) def listen_channels(self): - self.channel_tasks.append(create_task(self._recv_iopub(), self.task_group)) - self.channel_tasks.append(create_task(self._recv_shell(), self.task_group)) + self.channel_tasks.append(self.task_group.create_task(self._recv_iopub())) + self.channel_tasks.append(self.task_group.create_task(self._recv_shell())) async def stop(self) -> None: self.kernel_process.kill() @@ -97,13 +97,13 @@ async def _recv_iopub(self): while True: msg = await self.receive_message(self.iopub_channel, change_str_to_date=True) msg["channel"] = "iopub" - self.recv_queue.put_nowait(msg) + self.recv_queue.send_nowait(msg) async def _recv_shell(self): while True: msg = await self.receive_message(self.shell_channel, change_str_to_date=True) msg["channel"] = "shell" - self.recv_queue.put_nowait(msg) + self.recv_queue.send_nowait(msg) async def send_message( self, diff --git a/plugins/local_terminals/pyproject.toml b/plugins/local_terminals/pyproject.toml index 3c81593..c6f0cf8 100644 --- a/plugins/local_terminals/pyproject.toml +++ b/plugins/local_terminals/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "pywinpty;platform_system=='Windows'", "anyio >=3.7.0,<5",] +dependencies = [ "txl ==0.3.3", "pywinpty;platform_system=='Windows'", "anyio >=4.14.2,<5",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/local_terminals/txl_local_terminals/terminal.py b/plugins/local_terminals/txl_local_terminals/terminal.py index 513d6eb..e3fa1cf 100644 --- a/plugins/local_terminals/txl_local_terminals/terminal.py +++ b/plugins/local_terminals/txl_local_terminals/terminal.py @@ -1,16 +1,18 @@ import fcntl +import math import os import pty import shlex import struct import termios +from typing import Any -from anyio import wait_readable -from anyioutils import Event, Queue +from anyio import Event, create_memory_object_stream, wait_readable from textual.widget import Widget from textual.widgets._header import HeaderTitle from txl.base import Header, TerminalFactory, Terminals +from txl.stapled import StapledObjectStream class TerminalsMeta(type(Terminals), type(Widget)): @@ -22,8 +24,12 @@ def __init__(self, task_group, header: Header, terminal: TerminalFactory): self.task_group = task_group self.header = header self.terminal = terminal - self._send_queue = Queue() - self._recv_queue = Queue() + self._send_queue = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) + self._recv_queue = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) self._data_or_disconnect = None self._event = Event() super().__init__() @@ -66,9 +72,9 @@ async def _receive(self): self._event.set() async def _run(self): - await self._send_queue.put(["setup", {}]) + await self._send_queue.send(["setup", {}]) while True: - msg = await self._recv_queue.get() + msg = await self._recv_queue.receive() if msg[0] == "stdin" and msg[1] is not None: self._p_out.write(msg[1].encode()) elif msg[0] == "set_size": @@ -78,8 +84,8 @@ async def _run(self): async def _send(self): while True: await self._event.wait() - self._event.clear() + self._event = Event() if self._data_or_disconnect is None: - await self._send_queue.put(["disconnect", 1]) + await self._send_queue.send(["disconnect", 1]) else: - await self._send_queue.put(["stdout", self._data_or_disconnect]) + await self._send_queue.send(["stdout", self._data_or_disconnect]) diff --git a/plugins/local_terminals/txl_local_terminals/win_terminal.py b/plugins/local_terminals/txl_local_terminals/win_terminal.py index 7ee1c8f..a163700 100644 --- a/plugins/local_terminals/txl_local_terminals/win_terminal.py +++ b/plugins/local_terminals/txl_local_terminals/win_terminal.py @@ -1,13 +1,15 @@ +import math import os from functools import partial +from typing import Any -from anyio import create_task_group, to_thread -from anyioutils import Event, Queue +from anyio import create_memory_object_stream, create_task_group, to_thread from textual.widget import Widget from textual.widgets._header import HeaderTitle from winpty import PTY from txl.base import Header, TerminalFactory, Terminals +from txl.stapled import StapledObjectStream class TerminalsMeta(type(Terminals), type(Widget)): @@ -19,10 +21,13 @@ def __init__(self, task_group, header: Header, terminal: TerminalFactory): self.task_group = task_group self.header = header self.terminal = terminal - self._send_queue = Queue() - self._recv_queue = Queue() + self._send_queue = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) + self._recv_queue = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) self._data_or_disconnect = None - self._event = Event() super().__init__() async def open(self): @@ -44,7 +49,7 @@ def _open_terminal(self): return process async def _run(self): - await self._send_queue.put(["setup", {}]) + await self._send_queue.send(["setup", {}]) async with create_task_group() as tg: tg.start_soon(self._send) @@ -55,15 +60,15 @@ async def _send(self): try: data = await to_thread.run_sync(partial(self._process.read, blocking=True)) except Exception: - await self._send_queue.put(["disconnect", 1]) + await self._send_queue.send(["disconnect", 1]) return else: - await self._send_queue.put(["stdout", data]) + await self._send_queue.send(["stdout", data]) async def _recv(self): while True: try: - msg = await self._recv_queue.get() + msg = await self._recv_queue.receive() except Exception: return if msg[0] == "stdin": diff --git a/plugins/markdown_viewer/txl_markdown_viewer/main.py b/plugins/markdown_viewer/txl_markdown_viewer/main.py index a616dbe..aa452d8 100644 --- a/plugins/markdown_viewer/txl_markdown_viewer/main.py +++ b/plugins/markdown_viewer/txl_markdown_viewer/main.py @@ -1,5 +1,4 @@ from anyio import create_task_group, sleep -from anyioutils import create_task from fps import Module from textual._context import active_app from textual.app import App @@ -34,7 +33,7 @@ async def _update_viewer(self): self.mount(self.viewer) def on_change(self, target, event): - create_task(self.update_viewer(), self.task_group) + self.task_group.create_task(self.update_viewer()) class MarkdownViewerModule(Module): diff --git a/plugins/notebook_editor/pyproject.toml b/plugins/notebook_editor/pyproject.toml index 8b9979c..b6bce93 100644 --- a/plugins/notebook_editor/pyproject.toml +++ b/plugins/notebook_editor/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "pycrdt >=0.12.44,<0.13.0", "jupyter-ydoc >=3.0.2,<4.0.0",] +dependencies = [ "txl ==0.3.3", "pycrdt >=0.14.1,<0.15.0", "jupyter-ydoc >=4.1.1,<5.0.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/notebook_editor/txl_notebook_editor/main.py b/plugins/notebook_editor/txl_notebook_editor/main.py index 84f5987..dc23689 100644 --- a/plugins/notebook_editor/txl_notebook_editor/main.py +++ b/plugins/notebook_editor/txl_notebook_editor/main.py @@ -1,9 +1,9 @@ import json +import math from importlib.metadata import entry_points from typing import Any import anyio -from anyioutils import Queue, TaskGroup from fps import Module from httpx import AsyncClient from textual._context import active_app @@ -25,6 +25,7 @@ Launcher, MainArea, ) +from txl.stapled import StapledObjectStream ydocs = {ep.name: ep.load() for ep in entry_points(group="jupyter_ydoc")} @@ -84,8 +85,12 @@ def __init__( self.cell_i = 0 self.cell_copy = None self.edit_mode = False - self.nb_change_target = Queue() - self.nb_change_events = Queue() + self.nb_change_target = StapledObjectStream( + *anyio.create_memory_object_stream[Any](max_buffer_size=math.inf) + ) + self.nb_change_events = StapledObjectStream( + *anyio.create_memory_object_stream[Any](max_buffer_size=math.inf) + ) self.top_bar = TopBar() def compose(self) -> ComposeResult: @@ -169,13 +174,13 @@ def update(self): self.cells[self.cell_i].select() def on_change(self, target, events): - self.nb_change_target.put_nowait(target) - self.nb_change_events.put_nowait(events) + self.nb_change_target.send_nowait(target) + self.nb_change_events.send_nowait(events) async def observe_nb_changes(self): while True: - target = await self.nb_change_target.get() - events = await self.nb_change_events.get() + target = await self.nb_change_target.receive() + events = await self.nb_change_events.receive() if target == "meta": for event in events: meta = event.target @@ -413,7 +418,7 @@ async def start(self) -> None: _kernelspecs = await kernelspecs.get() - async with TaskGroup() as self.tg: + async with anyio.create_task_group() as self.tg: def notebook_editor_factory(): active_app.set(app) return NotebookEditor( diff --git a/plugins/remote_contents/pyproject.toml b/plugins/remote_contents/pyproject.toml index bb16a7e..baaaf3e 100644 --- a/plugins/remote_contents/pyproject.toml +++ b/plugins/remote_contents/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "httpx>=0.23.1", "httpx-ws>=0.4.2", "pycrdt >=0.12.44,<0.13.0", "pycrdt-websocket >=0.16.0,<0.17.0", "jupyter-ydoc >=3.0.2,<4.0.0",] +dependencies = [ "txl ==0.3.3", "httpx>=0.23.1", "httpx-ws>=0.4.2", "pycrdt >=0.14.1,<0.15.0", "pycrdt-websocket >=0.16.4,<0.17.0", "jupyter-ydoc >=4.1.1,<5.0.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/remote_kernels/pyproject.toml b/plugins/remote_kernels/pyproject.toml index f20c126..99a331d 100644 --- a/plugins/remote_kernels/pyproject.toml +++ b/plugins/remote_kernels/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "txl_kernel", "httpx >=0.23.1", "httpx-ws >=0.4.2", "pycrdt >=0.12.44,<0.13.0",] +dependencies = [ "txl ==0.3.3", "txl_kernel", "httpx >=0.23.1", "httpx-ws >=0.4.2", "pycrdt >=0.14.1,<0.15.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/plugins/remote_kernels/txl_remote_kernels/driver.py b/plugins/remote_kernels/txl_remote_kernels/driver.py index fc2ee4c..18be423 100644 --- a/plugins/remote_kernels/txl_remote_kernels/driver.py +++ b/plugins/remote_kernels/txl_remote_kernels/driver.py @@ -6,7 +6,6 @@ import httpx from anyio import Lock, sleep -from anyioutils import create_task from httpx_ws import aconnect_ws from txl_kernel.driver import KernelMixin from txl_kernel.message import date_to_str @@ -40,7 +39,7 @@ def __init__( self.cookies = httpx.Cookies() i = self.base_url.find(":") self.ws_url = ("wss" if self.base_url[i - 1] == "s" else "ws") + self.base_url[i:] - self.start_task = create_task(self.start(), task_group) + self.start_task = task_group.create_task(self.start()) self.comm_handlers = comm_handlers self.shell_channel = "shell" self.control_channel = "control" @@ -79,7 +78,7 @@ async def start(self): cookies=self.cookies, subprotocols=["v1.kernel.websocket.jupyter.org"], ) as self.websocket: - recv_task = create_task(self._recv(), self.task_group) + recv_task = self.task_group.create_task(self._recv()) try: await self.wait_for_ready() self.started.set() @@ -98,7 +97,7 @@ async def _recv(self): msg = json.loads(message.data) else: msg = from_binary(message.data) - self.recv_queue.put_nowait(msg) + self.recv_queue.send_nowait(msg) async def send_message( self, diff --git a/plugins/remote_terminals/txl_remote_terminals/main.py b/plugins/remote_terminals/txl_remote_terminals/main.py index 28adc93..c8edc8f 100644 --- a/plugins/remote_terminals/txl_remote_terminals/main.py +++ b/plugins/remote_terminals/txl_remote_terminals/main.py @@ -1,15 +1,16 @@ -from typing import Dict, List +import math +from typing import Any, Dict, List from urllib import parse import httpx -from anyio import create_task_group, sleep -from anyioutils import Event, Queue, create_task +from anyio import Event, create_memory_object_stream, create_task_group, sleep from fps import Module from httpx_ws import aconnect_ws from textual.widget import Widget from textual.widgets._header import HeaderTitle from txl.base import Header, Launcher, TerminalFactory, Terminals +from txl.stapled import StapledObjectStream class TerminalsMeta(type(Terminals), type(Widget)): @@ -34,8 +35,12 @@ def __init__( self.task_group = task_group i = base_url.find(":") self.ws_url = ("wss" if base_url[i - 1] == "s" else "ws") + base_url[i:] - self._recv_queue = Queue() - self._send_queue = Queue() + self._recv_queue = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) + self._send_queue = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) self._done = Event() super().__init__() @@ -65,12 +70,12 @@ async def open(self): f"{self.ws_url}/terminals/websocket/{name}", cookies=self.cookies ) as self.websocket: self.task_group.start_soon(self._recv) - self.send_task = create_task(self._send(), self.task_group) + self.send_task = self.task_group.create_task(self._send()) await self._done.wait() async def _send(self): while True: - message = await self._send_queue.get() + message = await self._send_queue.receive() try: await self.websocket.send_json(message) except BaseException: @@ -84,7 +89,7 @@ async def _recv(self): except Exception: self.send_task.cancel() return - await self._recv_queue.put(message) + await self._recv_queue.send(message) class RemoteTerminalsModule(Module): diff --git a/plugins/terminal/txl_terminal/main.py b/plugins/terminal/txl_terminal/main.py index 0d60624..69e7f70 100644 --- a/plugins/terminal/txl_terminal/main.py +++ b/plugins/terminal/txl_terminal/main.py @@ -1,8 +1,7 @@ from functools import partial import pyte -from anyio import create_task_group, sleep -from anyioutils import Event, create_task +from anyio import Event, create_task_group, sleep from fps import Module from rich.console import RenderableType from rich.text import Text @@ -45,7 +44,7 @@ def __init__(self, send_queue, recv_queue, main_area: MainArea, task_group): self._recv_queue = recv_queue self._display = PyteDisplay([Text()]) self.size_set = Event() - create_task(self._recv(), task_group) + task_group.create_task(self._recv()) main_area.set_label("Terminal") def render(self) -> RenderableType: @@ -60,16 +59,16 @@ def on_resize(self, event: events.Resize): async def on_key(self, event: events.Key) -> None: char = CTRL_KEYS.get(event.key) or event.character - await self._send_queue.put(["stdin", char]) + await self._send_queue.send(["stdin", char]) event.stop() async def _recv(self): await self.size_set.wait() while True: - message = await self._recv_queue.get() + message = await self._recv_queue.receive() cmd = message[0] if cmd == "setup": - await self._send_queue.put(["set_size", self._nrow, self._ncol, 567, 573]) + await self._send_queue.send(["set_size", self._nrow, self._ncol, 567, 573]) elif cmd == "stdout": chars = message[1] self._stream.feed(chars) diff --git a/plugins/text_editor/txl_text_editor/main.py b/plugins/text_editor/txl_text_editor/main.py index c19a03b..e32f046 100644 --- a/plugins/text_editor/txl_text_editor/main.py +++ b/plugins/text_editor/txl_text_editor/main.py @@ -1,5 +1,7 @@ -from anyio import create_task_group, sleep -from anyioutils import Queue +import math +from typing import Any + +from anyio import create_memory_object_stream, create_task_group, sleep from fps import Module from textual._context import active_app from textual.app import App @@ -8,6 +10,7 @@ from textual.keys import Keys from txl.base import Contents, Editor, Editors, MainArea +from txl.stapled import StapledObjectStream from txl.text_input import TextInput @@ -24,8 +27,12 @@ def __init__(self, contents: Contents, main_area: MainArea, task_group) -> None: self.contents = contents self.main_area = main_area self.task_group = task_group - self.change_target = Queue() - self.change_events = Queue() + self.change_target = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) + self.change_events = StapledObjectStream( + *create_memory_object_stream[Any](max_buffer_size=math.inf) + ) async def open(self, path: str) -> None: self.path = path @@ -40,13 +47,13 @@ async def close(self) -> None: await self.editor.stop() def on_change(self, target, events): - self.change_target.put_nowait(target) - self.change_events.put_nowait(events) + self.change_target.send_nowait(target) + self.change_events.send_nowait(events) async def observe_changes(self): while True: - target = await self.change_target.get() - events = await self.change_events.get() + target = await self.change_target.receive() + events = await self.change_events.receive() if target == "state": if "dirty" in events.keys: dirty = events.keys["dirty"]["newValue"] diff --git a/plugins/widgets/pyproject.toml b/plugins/widgets/pyproject.toml index 535241b..e2f1b4d 100644 --- a/plugins/widgets/pyproject.toml +++ b/plugins/widgets/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "txl ==0.3.3", "ypywidgets >=0.9.4,<0.10.0", "ypywidgets-textual >=0.5.4,<0.6.0", "pycrdt >=0.12.44,<0.13.0",] +dependencies = [ "txl ==0.3.3", "ypywidgets >=0.9.9,<0.10.0", "ypywidgets-textual >=0.5.5,<0.6.0", "pycrdt >=0.14.1,<0.15.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/txl/pyproject.toml b/txl/pyproject.toml index 61bf825..c61275b 100644 --- a/txl/pyproject.toml +++ b/txl/pyproject.toml @@ -11,7 +11,7 @@ requires-python = ">=3.10" license = "MIT" keywords = [] classifiers = [ "Development Status :: 4 - Beta", "Programming Language :: Python", "Programming Language :: Python :: 3.10", "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", "Programming Language :: Python :: Implementation :: CPython", "Programming Language :: Python :: Implementation :: PyPy",] -dependencies = [ "fps >=0.5.1,<0.6.0", "textual[syntax] >=6.8.0",] +dependencies = [ "fps >=0.6.5,<0.7.0", "textual[syntax] >=6.8.0",] [[project.authors]] name = "David Brochart" email = "david.brochart@gmail.com" diff --git a/txl/txl/stapled.py b/txl/txl/stapled.py new file mode 100644 index 0000000..26a092e --- /dev/null +++ b/txl/txl/stapled.py @@ -0,0 +1,11 @@ +from typing import TypeVar + +from anyio.streams.stapled import StapledObjectStream as _StapledObjectStream + +T_Item = TypeVar("T_Item") + + +# FIXME: remove when https://github.com/agronholm/anyio/pull/1241 is released +class StapledObjectStream(_StapledObjectStream[T_Item]): + def send_nowait(self, item: T_Item) -> None: + self.send_stream.send_nowait(item) # type: ignore[attr-defined]