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
31 changes: 31 additions & 0 deletions docs/async-ecosystem-tests.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,37 @@ aiosqlite 0.22.1 and uvloop 0.23.0. These cases did not reproduce a new rsloop
defect; no production implementation change was needed. See the PR checks for
platform-specific execution results.

## Greenlet re-entry and transport ownership regressions

`tests/ecosystem/test_greenlet_reentry.py` covers the two-run lifecycle used by
ASGI servers: start background tasks, then re-enter the loop from a deeper C
stack. It tests a minimal greenlet bridge and SQLAlchemy's `greenlet_spawn` at
depths 0, 2, 4, 8, 50 and 100. Each case runs in a subprocess with a deadline and
faulthandler, verifies progress by both contending tasks, and verifies cancellation
completes. No database is required. The active ready queue lives on the heap so
greenlet stack copying cannot invalidate its thread-local pointer.

`tests/test_transport_protocol_release.py` checks that closed TCP and subprocess
transports release protocols that retain their transport, including when
`connection_lost` raises. It also checks release of `StreamReaderProtocol` and its
cached reader. These tests use weak references and garbage collection rather than
RSS, which can vary with allocator behavior. Stream teardown clears the protocol,
cached bound methods, fast-reader references and saved context after delivering
`connection_lost`; `get_protocol()` then returns `None`. This breaks the closed
connection cycle without making native transports cyclic-GC types.

Run both regression groups with:

```sh
uv run --no-sync python scripts/run_python_tests.py \
tests/test_transport_protocol_release.py tests/ecosystem/test_greenlet_reentry.py -m ''
```

These regressions reproduce the scheduling and ownership mechanisms reported in
[issue #115](https://github.com/RustedBytes/rsloop/issues/115) and
[issue #116](https://github.com/RustedBytes/rsloop/issues/116). They do not replace
end-to-end Granian/PostgreSQL load or shutdown testing.

## PostgreSQL contracts

`tests/ecosystem/test_postgresql.py` adds 36 cases: asyncpg and SQLAlchemy
Expand Down
16 changes: 10 additions & 6 deletions src/engine/loop_core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -723,8 +723,10 @@ impl LoopCore {
ensure_running_loop(py, &loop_obj)?;
self.mark_runtime_thread();
let mut local_timers = LocalTimers::new(self);
let mut local_ready = VecDeque::new();
self.install_local_ready_queue(&mut local_ready);
// Greenlets can copy this run's C stack out while a callback executes.
// TLS must point to heap storage, not a queue header on that stack.
let mut local_ready = Box::new(VecDeque::new());
self.install_local_ready_queue(&mut *local_ready);

let mut pending_signal_error: Option<PyErr> = None;
let mut ready_batch = VecDeque::new();
Expand Down Expand Up @@ -787,7 +789,7 @@ impl LoopCore {
// hot stream of locally-scheduled callbacks.
if !local_ready.is_empty() {
if ready_batch.is_empty() {
std::mem::swap(&mut ready_batch, &mut local_ready);
std::mem::swap(&mut ready_batch, &mut *local_ready);
} else {
ready_batch.extend(local_ready.drain(..));
}
Expand Down Expand Up @@ -1686,9 +1688,11 @@ impl LoopCore {
return Err(item);
}

// SAFETY: `ready` points to run_forever's stack-local queue on this
// thread. Neither callback invocation nor runtime polling holds a
// mutable reference to it across this call.
// SAFETY: `ready` points to run_forever's heap-allocated queue on
// this thread, stable even across greenlet stack switches. Neither
// callback invocation nor runtime polling holds a mutable reference
// to it across this call. TLS is cleared before the allocation
// drops.
unsafe { (*ready).push_back(item) };
if !tls.drain_active.get() {
// I/O futures are polled with Python detached, but still on
Expand Down
16 changes: 14 additions & 2 deletions src/transport/process/core_protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -256,8 +256,20 @@ impl ProcessTransportCore {
let arg = exc
.map(|err| err.value(py).clone().unbind().into_any())
.unwrap_or_else(|| py.None());
self.call_protocol_method1(py, "connection_lost", arg)?;
Ok(())
let result = self.call_protocol_method1(py, "connection_lost", arg);
// Subprocess protocols commonly store their transport too. Release
// references even if the callback fails, outside the state lock so
// Python finalizers can safely reenter the transport.
let released = {
let mut state = self.state.lock().expect("poisoned process state");
state.context_needs_run = false;
(
std::mem::replace(&mut state.protocol, py.None()),
std::mem::replace(&mut state.context, py.None()),
)
};
drop(released);
result.map(|_| ())
}

#[cfg_attr(
Expand Down
27 changes: 22 additions & 5 deletions src/transport/stream/core_protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -485,14 +485,15 @@ impl StreamTransportCore {
)
};

if let Some(fast_path) = fast_path.as_ref() {
fast_path.connection_lost(py, exc)?;
let result = if let Some(fast_path) = fast_path.as_ref() {
fast_path.connection_lost(py, exc)
} else {
let arg = exc
.map(|err| err.value(py).clone().unbind().into_any())
.unwrap_or_else(|| py.None());
self.call_protocol_method1(py, &callback, &context, context_needs_run, arg)?;
}
self.call_protocol_method1(py, &callback, &context, context_needs_run, arg)
.map(|_| ())
};

// Preserve asyncio's post-close get_extra_info("socket") behavior
// without retaining the live shared owner through Python reference
Expand All @@ -504,6 +505,22 @@ impl StreamTransportCore {
if let Some(server) = server.and_then(|weak| weak.upgrade()) {
server.connection_lost();
}
Ok(())
// Even a failing connection_lost must break protocol -> transport ->
// protocol cycles. Bound methods and fast-reader caches also own Python
// references. Drop them outside the mutex: finalizers may reenter us.
let released = {
let mut state = self.state.lock().expect("poisoned transport state");
state.context_needs_run = false;
(
std::mem::replace(&mut state.protocol, py.None()),
std::mem::replace(
&mut state.callbacks,
super::protocol::ProtocolCallbacks::cleared(py),
),
std::mem::replace(&mut state.context, py.None()),
)
};
drop(released);
result
}
}
19 changes: 19 additions & 0 deletions src/transport/stream/protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,25 @@ pub(super) struct ProtocolCallbacks {
pub(super) stream_reader_fast_path: Option<StreamReaderFastPath>,
}

impl ProtocolCallbacks {
/// No protocol callbacks may run after connection_lost has been delivered.
/// Release cached bound methods and fast readers as well as the protocol.
#[cfg_attr(feature = "profile", hotpath::measure(impl_type = "ProtocolCallbacks"))]
pub(super) fn cleared(py: Python<'_>) -> Self {
Self {
connection_made: py.None(),
data_received: None,
eof_received: None,
connection_lost: py.None(),
pause_writing: py.None(),
resume_writing: py.None(),
get_buffer: None,
buffer_updated: None,
stream_reader_fast_path: None,
}
}
}

pub(super) enum StreamReaderFastPath {
Native {
protocol: Py<PyFastStreamProtocol>,
Expand Down
98 changes: 98 additions & 0 deletions tests/ecosystem/test_greenlet_reentry.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
"""Exercise greenlet stack copying in a subprocess: regressions can segfault."""

import asyncio
import subprocess
import sys
from pathlib import Path
from typing import cast

import pytest
import rsloop

pytestmark = pytest.mark.ecosystem


def _exercise(depth, bridge):
import greenlet

class AsyncGreenlet(greenlet.greenlet):
def __init__(self, fn, driver):
super().__init__(fn, driver)
self.driver = driver

def await_only(awaitable):
return cast(AsyncGreenlet, greenlet.getcurrent()).driver.switch(awaitable)

async def spawn(fn):
child = AsyncGreenlet(fn, greenlet.getcurrent())
result = child.switch()
while not child.dead:
try:
value = await result
except BaseException as exc: # noqa: BLE001 - forward cancellation into the greenlet
result = child.throw(type(exc), exc, exc.__traceback__)
else:
result = child.switch(value)
return result

if bridge == "sqlalchemy":
from sqlalchemy.util.concurrency import await_only, greenlet_spawn

spawn = greenlet_spawn

lock = asyncio.Lock()
progress = [0, 0]

def critical_section():
await_only(lock.acquire())
try:
await_only(asyncio.sleep(0.001))
finally:
lock.release()

async def worker(index):
while True:
await spawn(critical_section)
progress[index] += 1
await asyncio.sleep(0)

async def start():
return [asyncio.create_task(worker(i)) for i in range(2)]

def deeper(n, fn):
# map.__next__ deliberately adds C frames; a generator changes the repro.
return fn() if n == 0 else next(map(lambda _: deeper(n - 1, fn), [None])) # noqa: C417

loop = rsloop.new_event_loop()
asyncio.set_event_loop(loop)
tasks = loop.run_until_complete(start())
deeper(depth, lambda: loop.run_until_complete(asyncio.sleep(0.15)))
assert min(progress) > 1, progress
for task in tasks:
task.cancel()
loop.run_until_complete(
asyncio.wait_for(asyncio.gather(*tasks, return_exceptions=True), 2)
)
assert all(task.cancelled() for task in tasks)
loop.close()
asyncio.set_event_loop(None)


@pytest.mark.parametrize("depth", [0, 2, 4, 8, 50, 100])
@pytest.mark.parametrize("bridge", ["minimal", "sqlalchemy"])
def test_greenlet_wakeup_after_deeper_loop_reentry(depth, bridge):
pytest.importorskip("greenlet")
if bridge == "sqlalchemy":
pytest.importorskip("sqlalchemy")
script = (
f"import runpy; runpy.run_path({str(Path(__file__))!r})"
f"['_exercise']({depth}, {bridge!r})"
)
result = subprocess.run(
[sys.executable, "-X", "faulthandler", "-c", script],
capture_output=True,
text=True,
timeout=10,
check=False,
)
assert result.returncode == 0, result.stdout + result.stderr
Loading
Loading