mirror of
https://github.com/ntex-rs/ntex.git
synced 2025-04-04 13:27:39 +03:00
Handle socket close for poll driver (#549)
This commit is contained in:
parent
e4f24ee41f
commit
f5ee55d598
5 changed files with 35 additions and 36 deletions
|
@ -1,5 +1,9 @@
|
||||||
# Changes
|
# Changes
|
||||||
|
|
||||||
|
## [2.5.10] - 2025-03-28
|
||||||
|
|
||||||
|
* Better closed sockets handling
|
||||||
|
|
||||||
## [2.5.9] - 2025-03-27
|
## [2.5.9] - 2025-03-27
|
||||||
|
|
||||||
* Handle closed sockets
|
* Handle closed sockets
|
||||||
|
|
|
@ -40,7 +40,7 @@ ntex-util = "2.5"
|
||||||
|
|
||||||
ntex-tokio = { version = "0.5.3", optional = true }
|
ntex-tokio = { version = "0.5.3", optional = true }
|
||||||
ntex-compio = { version = "0.2.4", optional = true }
|
ntex-compio = { version = "0.2.4", optional = true }
|
||||||
ntex-neon = { version = "0.1.14", optional = true }
|
ntex-neon = { version = "0.1.15", optional = true }
|
||||||
|
|
||||||
bitflags = { workspace = true }
|
bitflags = { workspace = true }
|
||||||
cfg-if = { workspace = true }
|
cfg-if = { workspace = true }
|
||||||
|
|
|
@ -1,5 +1,5 @@
|
||||||
use std::os::fd::{AsRawFd, RawFd};
|
use std::os::fd::{AsRawFd, RawFd};
|
||||||
use std::{cell::Cell, cell::RefCell, future::Future, io, rc::Rc, task, task::Poll};
|
use std::{cell::Cell, cell::RefCell, future::Future, io, mem, rc::Rc, task, task::Poll};
|
||||||
|
|
||||||
use ntex_neon::driver::{DriverApi, Event, Handler};
|
use ntex_neon::driver::{DriverApi, Event, Handler};
|
||||||
use ntex_neon::{syscall, Runtime};
|
use ntex_neon::{syscall, Runtime};
|
||||||
|
@ -18,7 +18,6 @@ bitflags::bitflags! {
|
||||||
struct Flags: u8 {
|
struct Flags: u8 {
|
||||||
const RD = 0b0000_0001;
|
const RD = 0b0000_0001;
|
||||||
const WR = 0b0000_0010;
|
const WR = 0b0000_0010;
|
||||||
const CLOSED = 0b0000_0100;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -106,18 +105,15 @@ impl<T> Handler for StreamOpsHandler<T> {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
let item = &mut streams[id];
|
let item = &mut streams[id];
|
||||||
log::debug!("{}: FD event {:?} event: {:?}", item.tag(), id, ev);
|
if item.io.is_none() {
|
||||||
|
|
||||||
if item.flags.contains(Flags::CLOSED) {
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
log::debug!("{}: FD event {:?} event: {:?}", item.tag(), id, ev);
|
||||||
|
|
||||||
// handle HUP
|
// handle HUP
|
||||||
if ev.is_interrupt() {
|
if ev.is_interrupt() {
|
||||||
item.context.stopped(None);
|
item.context.stopped(None);
|
||||||
if item.io.take().is_some() {
|
close(id as u32, item, &self.inner.api, None, true);
|
||||||
close(id as u32, item, &self.inner.api);
|
|
||||||
}
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -177,9 +173,7 @@ impl<T> Handler for StreamOpsHandler<T> {
|
||||||
item.fd,
|
item.fd,
|
||||||
item.io.is_some()
|
item.io.is_some()
|
||||||
);
|
);
|
||||||
if item.io.is_some() {
|
close(id, &mut item, &self.inner.api, None, true);
|
||||||
close(id, &mut item, &self.inner.api);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
self.inner.delayd_drop.set(false);
|
self.inner.delayd_drop.set(false);
|
||||||
|
@ -197,10 +191,7 @@ impl<T> Handler for StreamOpsHandler<T> {
|
||||||
item.fd,
|
item.fd,
|
||||||
err
|
err
|
||||||
);
|
);
|
||||||
item.context.stopped(Some(err));
|
close(id as u32, item, &self.inner.api, Some(err), false);
|
||||||
if item.io.take().is_some() {
|
|
||||||
close(id as u32, item, &self.inner.api);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
@ -222,14 +213,26 @@ fn close<T>(
|
||||||
id: u32,
|
id: u32,
|
||||||
item: &mut StreamItem<T>,
|
item: &mut StreamItem<T>,
|
||||||
api: &DriverApi,
|
api: &DriverApi,
|
||||||
) -> ntex_rt::JoinHandle<io::Result<i32>> {
|
error: Option<io::Error>,
|
||||||
|
shutdown: bool,
|
||||||
|
) -> Option<ntex_rt::JoinHandle<io::Result<i32>>> {
|
||||||
|
if let Some(io) = item.io.take() {
|
||||||
|
log::debug!("{}: Closing ({}), {:?}", item.tag(), id, item.fd);
|
||||||
|
mem::forget(io);
|
||||||
|
if let Some(err) = error {
|
||||||
|
item.context.stopped(Some(err));
|
||||||
|
}
|
||||||
let fd = item.fd;
|
let fd = item.fd;
|
||||||
item.flags.insert(Flags::CLOSED);
|
|
||||||
api.detach(fd, id);
|
api.detach(fd, id);
|
||||||
ntex_rt::spawn_blocking(move || {
|
Some(ntex_rt::spawn_blocking(move || {
|
||||||
syscall!(libc::shutdown(fd, libc::SHUT_RDWR))?;
|
if shutdown {
|
||||||
|
let _ = syscall!(libc::shutdown(fd, libc::SHUT_RDWR));
|
||||||
|
}
|
||||||
syscall!(libc::close(fd))
|
syscall!(libc::close(fd))
|
||||||
})
|
}))
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<T> StreamCtl<T> {
|
impl<T> StreamCtl<T> {
|
||||||
|
@ -237,13 +240,7 @@ impl<T> StreamCtl<T> {
|
||||||
let id = self.id as usize;
|
let id = self.id as usize;
|
||||||
let fut = self.inner.with(|streams| {
|
let fut = self.inner.with(|streams| {
|
||||||
let item = &mut streams[id];
|
let item = &mut streams[id];
|
||||||
if let Some(io) = item.io.take() {
|
close(self.id, item, &self.inner.api, None, false)
|
||||||
log::debug!("{}: Closing ({}), {:?}", item.tag(), id, item.fd);
|
|
||||||
std::mem::forget(io);
|
|
||||||
Some(close(self.id, item, &self.inner.api))
|
|
||||||
} else {
|
|
||||||
None
|
|
||||||
}
|
|
||||||
});
|
});
|
||||||
async move {
|
async move {
|
||||||
if let Some(fut) = fut {
|
if let Some(fut) = fut {
|
||||||
|
@ -360,9 +357,7 @@ impl<T> Drop for StreamCtl<T> {
|
||||||
item.fd,
|
item.fd,
|
||||||
item.io.is_some()
|
item.io.is_some()
|
||||||
);
|
);
|
||||||
if item.io.is_some() {
|
close(self.id, &mut item, &self.inner.api, None, true);
|
||||||
close(self.id, &mut item, &self.inner.api);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
self.inner.streams.set(Some(streams));
|
self.inner.streams.set(Some(streams));
|
||||||
} else {
|
} else {
|
||||||
|
|
|
@ -1,6 +1,6 @@
|
||||||
[package]
|
[package]
|
||||||
name = "ntex-server"
|
name = "ntex-server"
|
||||||
version = "2.7.3"
|
version = "2.7.4"
|
||||||
authors = ["ntex contributors <team@ntex.rs>"]
|
authors = ["ntex contributors <team@ntex.rs>"]
|
||||||
description = "Server for ntex framework"
|
description = "Server for ntex framework"
|
||||||
keywords = ["network", "framework", "async", "futures"]
|
keywords = ["network", "framework", "async", "futures"]
|
||||||
|
|
|
@ -180,7 +180,7 @@ impl<F: ServerConfiguration> HandleCmdState<F> {
|
||||||
fn process(&mut self, mut item: F::Item) {
|
fn process(&mut self, mut item: F::Item) {
|
||||||
loop {
|
loop {
|
||||||
if !self.workers.is_empty() {
|
if !self.workers.is_empty() {
|
||||||
if self.next > self.workers.len() {
|
if self.next >= self.workers.len() {
|
||||||
self.next = self.workers.len() - 1;
|
self.next = self.workers.len() - 1;
|
||||||
}
|
}
|
||||||
match self.workers[self.next].send(item) {
|
match self.workers[self.next].send(item) {
|
||||||
|
|
Loading…
Add table
Add a link
Reference in a new issue