FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

Release interpreter channel payloads outside internal locks by 1ndahous3 · Pull Request #8964 · RustPython/RustPython · GitHub

Repository navigation

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

Filter by extension

Filter by extension .py  (1) .rs  (1) All 2 file types selected
Viewed files
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Unified
Split
Hide whitespace
Diff view
Unified
Split
Hide whitespace
78 changes: 51 additions & 27 deletions crates/vm/src/stdlib/_interpchannels.rs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -366,18 +366,23 @@ pub(crate) mod _interpchannels {

/// `channel_destroy`: forget the channel entirely.
fn channel_destroy(cid: i64) -> ChanResult<Vec<(Arc<Waiting>, bool)>> {
let mut table = channels().lock();
let entry = table.refs.remove(&cid).ok_or(ChanErr::NotFound)?;
Ok(match entry.chan {
Some(chan) => drain_queue(&mut chan.state.lock()),
None => Vec::new(),
})
let entry = channels()
.lock()
.refs
.remove(&cid)
.ok_or(ChanErr::NotFound)?;
let items = match entry.chan {
Some(chan) => core::mem::take(&mut chan.state.lock().queue),
None => VecDeque::new(),
};
Ok(retire_items(items))
}

fn drain_queue(state: &mut ChannelState) -> Vec<(Arc<Waiting>, bool)> {
state
.queue
.drain(..)
// Shared buffers can run owner callbacks when released. Call only after
// unlocking both the directory and channel state.
fn retire_items(items: VecDeque<ChannelItem>) -> Vec<(Arc<Waiting>, bool)> {
items
.into_iter()
.filter_map(|item| item.waiting.map(|w| (w, false)))
.collect()
}
Expand All @@ -394,10 +399,12 @@ pub(crate) mod _interpchannels {
return Vec::new();
}
let entry = table.refs.remove(&cid).expect("just looked it up");
match entry.chan {
Some(chan) => drain_queue(&mut chan.state.lock()),
None => Vec::new(),
}
drop(table);
let items = match entry.chan {
Some(chan) => core::mem::take(&mut chan.state.lock().queue),
None => VecDeque::new(),
};
retire_items(items)
}

/// `_channels_add_id_object`.
Expand Down Expand Up @@ -456,12 +463,13 @@ pub(crate) mod _interpchannels {
state.closing = true;
return Ok(Vec::new());
}
let waiters = drain_queue(&mut state);
let items = core::mem::take(&mut state.queue);
state.open = false;
state.ends.release_all();
drop(state);
entry.chan = None;
Ok(waiters)
drop(table);
Ok(retire_items(items))
}

/// `channel_release`: close one or both ends for the current interpreter.
Expand All @@ -486,6 +494,7 @@ pub(crate) mod _interpchannels {
.filter_map(|r| r.chan.clone())
.collect();
let mut waiters = Vec::new();
let mut dropped = Vec::new();
for chan in chans {
let mut state = chan.state.lock();
let mut i = 0;
Expand All @@ -500,17 +509,19 @@ pub(crate) mod _interpchannels {
if let Some(w) = item.waiting {
waiters.push((w, false));
}
dropped.push(item.data);
continue;
}
// UNBOUND_ERROR / UNBOUND_REPLACE keep the slot but throw the
// data away; the item is now "unbound".
item.data = None;
dropped.push(item.data.take());
i += 1;
}
state.ends.clear_interpreter(interpid);
state.open = state.ends.is_open();
}
release_waiters_detached(waiters);
drop(dropped);
}

/// `resolve_unboundop`.
Expand Down Expand Up @@ -894,23 +905,33 @@ pub(crate) mod _interpchannels {
let cid = parse_cid(parsed[0].as_deref().unwrap(), vm)?.0;
let send = flag(parsed[1].as_deref(), vm)?;
let chan = channels_lookup(cid).map_err(|e| e.into_py(cid, vm))?;
let interpreters = crate::vm::runtime::list_interpreters();
let state = chan.state.lock();
if send && state.closing {
drop(state);
return Err(ChanErr::Closed.into_py(cid, vm));
}
let mut out = Vec::new();
for info in crate::vm::runtime::list_interpreters() {
for info in interpreters {
if state.ends.find(info.id, send).is_some_and(|i| {
if send {
state.ends.send[i].1
} else {
state.ends.recv[i].1
}
}) {
out.push(vm.ctx.new_int(info.id).into());
out.push(info.id);
}
}
Ok(vm.ctx.new_list(out).into())
drop(state);
Ok(vm
.ctx
.new_list(
out.into_iter()
.map(|id| vm.ctx.new_int(id).into())
.collect(),
)
.into())
}

/// `channel_send` up to its `closing` check, which precedes converting the
Expand Down Expand Up @@ -982,15 +1003,18 @@ pub(crate) mod _interpchannels {
if timed_out {
// The send is failing now, so make sure the item won't be received.
let mut state = chan.state.lock();
if let Some(pos) = state.queue.iter().position(|it| {
it.waiting
.as_ref()
.is_some_and(|w| Arc::ptr_eq(w, &waiting))
}) {
state.queue.remove(pos);
}
let removed = state
.queue
.iter()
.position(|it| {
it.waiting
.as_ref()
.is_some_and(|w| Arc::ptr_eq(w, &waiting))
})
.and_then(|pos| state.queue.remove(pos));
drop(state);
finish_closing(cid);
drop(removed);
if !waiting.state.lock().received {
return Err(vm
.new_os_subtype_error(
Expand Down
47 changes: 47 additions & 0 deletions extra_tests/snippets/stdlib_subinterpreters.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -117,3 +117,50 @@ def receive():
encoding="utf-8",
)
subprocess.run([sys.executable, str(script)], check=True, timeout=30)


# Releasing queued buffers can run a finalizer which reenters the channel API.
subprocess.run(
[
sys.executable,
"-c",
"""
import _interpchannels as channels
import _interpreters
import sys

released = []

class Exporter(bytearray):
def __del__(self):
channels.list_all()
released.append(True)

for operation in ("destroy", "close", "drop", "timeout"):
# CPython 3.14.6 also deadlocks when close/drop releases a reentrant exporter.
if sys.implementation.name == "cpython" and operation in ("close", "drop"):
continue
channel = channels.create(3)
released.clear()
if operation == "timeout":
try:
channels.send_buffer(channel, Exporter(b"data"), timeout=0)
except TimeoutError:
pass
else:
raise AssertionError("send must time out without a receiver")
else:
channels.send_buffer(channel, Exporter(b"data"), blocking=False)
assert not released
if operation == "destroy":
channels.destroy(channel)
elif operation == "close":
channels.close(channel, force=True)
else:
del channel
assert released == [True], operation
""",
],
check=True,
timeout=30,
)
Loading

Back | FazBrowse Home | New Git URL