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
22 changes: 8 additions & 14 deletions kernel/src/driver/net/veth.rs
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,12 @@ impl VethDriver {
let Some(iface) = self.inner.lock().self_iface_ref.upgrade() else {
return IngressDisposition::Local;
};
// Like Linux's ptype_all taps, packet sockets observe ingress before
// bridge/routing can consume it, not only frames delivered locally.
let packet_iface: Arc<dyn Iface> = iface.clone();
let pkt_type = crate::net::socket::packet::classify_packet(data, &packet_iface);
crate::net::socket::packet::deliver_to_packet_sockets(&packet_iface, data, pkt_type);

if let Some(bridge_data) = iface.common_bridge_data() {
Veth::to_bridge(&bridge_data, data);
return IngressDisposition::Consumed;
Expand Down Expand Up @@ -329,23 +335,14 @@ impl phy::TxToken for VethTxToken {

pub struct VethRxToken {
buffer: Vec<u8>,
driver: VethDriver,
}

impl RxToken for VethRxToken {
fn consume<R, F>(self, f: F) -> R
where
F: FnOnce(&[u8]) -> R,
{
let packet = self.buffer.as_slice();

// 向注册的 packet socket 分发数据包
if let Some(iface) = self.driver.iface() {
let pkt_type = crate::net::socket::packet::classify_packet(packet, &iface);
crate::net::socket::packet::deliver_to_packet_sockets(&iface, packet, pkt_type);
}

f(packet)
f(self.buffer.as_slice())
}
}

Expand All @@ -368,10 +365,7 @@ impl phy::Device for VethDriver {
guard.recv_local().map(|buf| {
// log::info!("VethDriver received data: {:?}", buf);
(
VethRxToken {
buffer: buf,
driver: self.clone(),
},
VethRxToken { buffer: buf },
VethTxToken {
driver: self.clone(),
},
Expand Down
5 changes: 5 additions & 0 deletions kernel/src/filesystem/vfs/file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2616,6 +2616,11 @@ impl File {
new_flags = FileFlags::from_bits_truncate(new_bits);

self.private_data.lock().update_flags(new_flags)?;
// Socket operations such as accept consult their cached mode. Keep it
// coherent for both F_SETFL and FIONBIO, under the same update lock.
if let Some(socket) = self.inode.as_socket() {
socket.set_nonblocking(new_flags.contains(FileFlags::O_NONBLOCK));
}
// 更新文件的打开模式
*self.flags.write() = new_flags;

Expand Down
10 changes: 0 additions & 10 deletions kernel/src/filesystem/vfs/syscall/sys_fcntl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -188,16 +188,6 @@ impl SysFcntlHandle {
}
}

// Keep socket object nonblocking state in sync with file flags.
// Some socket implementations consult an internal AtomicBool rather than
// the FileFlags, so fcntl(F_SETFL,O_NONBLOCK) must be propagated.
if file.file_type() == FileType::Socket {
if let Ok(inode) = ProcessManager::current_pcb().get_socket_inode(fd) {
if let Some(sock) = inode.as_socket() {
sock.set_nonblocking(new_flags.contains(FileFlags::O_NONBLOCK));
}
}
}
return Ok(0);
}

Expand Down
53 changes: 42 additions & 11 deletions kernel/src/net/socket/inet/stream/inner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -768,6 +768,23 @@ pub struct Listening {
}

impl Listening {
/// Ordinary passive opens become acceptable only after the final ACK.
/// A peer may already have sent FIN, so CLOSE_WAIT remains acceptable.
/// Snapshot both endpoints under the caller's SocketSet lock: a later RST
/// can clear the tuple before the accepted socket is constructed.
fn accept_endpoints(
socket: &tcp::Socket,
) -> Option<(smoltcp::wire::IpEndpoint, smoltcp::wire::IpEndpoint)> {
if matches!(
socket.state(),
tcp::State::Established | tcp::State::CloseWait
) {
Some((socket.local_endpoint()?, socket.remote_endpoint()?))
} else {
None
}
}

fn slot_capacity(backlog: usize) -> usize {
// Bounded per-interface emulation, not Linux's global accept queue.
// At 256 KiB per socket this uses at most 2 MiB per interface.
Expand Down Expand Up @@ -884,24 +901,26 @@ impl Listening {
pub fn accept(&mut self) -> Result<(Established, smoltcp::wire::IpEndpoint), SystemError> {
// Resizing can invalidate vector indices. Select the current live slot
// under the caller's inner write lock instead of caching a poll index.
let index = self
let (index, local_endpoint, remote_endpoint) = self
.inners
.iter()
.position(|bound| bound.with::<tcp::Socket, _, _>(|socket| socket.is_active()))
.enumerate()
.find_map(|(index, bound)| {
bound.with::<tcp::Socket, _, _>(|socket| {
Self::accept_endpoints(socket).map(|(local, peer)| (index, local, peer))
})
})
.ok_or(SystemError::EAGAIN_OR_EWOULDBLOCK)?;

let retire = self.can_remove_slot(index);
let connected = &mut self.inners[index];

let remote_endpoint = connected.with::<smoltcp::socket::tcp::Socket, _, _>(|socket| {
socket
.remote_endpoint()
.expect("A Connected Tcp With No Remote Endpoint")
});

if retire {
let connected = self.inners.remove(index);
return Ok((Established::new(connected, None), remote_endpoint));
return Ok((
Established::with_endpoints(connected, None, local_endpoint, remote_endpoint),
remote_endpoint,
));
}

// log::debug!("local at {:?}", local_endpoint);
Expand Down Expand Up @@ -932,7 +951,10 @@ impl Listening {
// TODO is smoltcp socket swappable?
core::mem::swap(&mut new_listen, connected);

return Ok((Established::new(new_listen, None), remote_endpoint));
return Ok((
Established::with_endpoints(new_listen, None, local_endpoint, remote_endpoint),
remote_endpoint,
));
}

pub fn update_io_events(&self, pollee: &AtomicUsize) {
Expand All @@ -956,7 +978,7 @@ impl Listening {

// log::info!("Listening::update_io_events");
let ready = self.inners.iter().any(|inner| {
inner.with::<smoltcp::socket::tcp::Socket, _, _>(|socket| socket.is_active())
inner.with::<tcp::Socket, _, _>(|socket| Self::accept_endpoints(socket).is_some())
});

if ready {
Expand Down Expand Up @@ -1015,6 +1037,15 @@ impl Established {
smoltcp::wire::IpAddress::Ipv4(smoltcp::wire::Ipv4Address::UNSPECIFIED),
0,
));
Self::with_endpoints(inner, reservation, local, peer)
}

fn with_endpoints(
inner: socket::inet::BoundInner,
reservation: Option<TcpPortReservation>,
local: smoltcp::wire::IpEndpoint,
peer: smoltcp::wire::IpEndpoint,
) -> Self {
Self {
inner,
local,
Expand Down
7 changes: 6 additions & 1 deletion kernel/src/net/socket/inode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,10 +30,15 @@ impl<T: Socket + 'static> IndexNode for T {
fn open(
&self,
data: MutexGuard<FilePrivateData>,
_: &crate::filesystem::vfs::file::FileFlags,
flags: &crate::filesystem::vfs::file::FileFlags,
) -> Result<(), SystemError> {
match &*data {
FilePrivateData::SocketCreate => {
// The new open file description owns the mode. In particular,
// accept does not inherit the listener's O_NONBLOCK.
self.set_nonblocking(
flags.contains(crate::filesystem::vfs::file::FileFlags::O_NONBLOCK),
);
self.open_file_counter().fetch_add(1, Ordering::Release);
Ok(())
}
Expand Down
6 changes: 5 additions & 1 deletion kernel/src/net/syscall/sys_socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -107,7 +107,11 @@ pub(super) fn do_socket(
is_close_on_exec,
)?;

let file = File::new_socket(inode, FileFlags::O_RDWR)?;
let mut file_flags = FileFlags::O_RDWR;
if is_nonblock {
file_flags.insert(FileFlags::O_NONBLOCK);
}
let file = File::new_socket(inode, file_flags)?;
// 把socket添加到当前进程的文件描述符表中
let current = ProcessManager::current_pcb();
current
Expand Down
8 changes: 6 additions & 2 deletions kernel/src/net/syscall/sys_socketpair.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,8 +176,12 @@ pub(super) fn do_socketpair(
}
};

let file_a = File::new_socket(socket_a, FileFlags::O_RDWR)?;
let file_b = File::new_socket(socket_b, FileFlags::O_RDWR)?;
let mut file_flags = FileFlags::O_RDWR;
if nonblocking {
file_flags.insert(FileFlags::O_NONBLOCK);
}
let file_a = File::new_socket(socket_a, file_flags)?;
let file_b = File::new_socket(socket_b, file_flags)?;
let file_a = alloc::sync::Arc::try_new(file_a).map_err(|_| SystemError::ENOMEM)?;
let file_b = alloc::sync::Arc::try_new(file_b).map_err(|_| SystemError::ENOMEM)?;
reservation.install_arc_pair(file_a, file_b)?;
Expand Down
2 changes: 2 additions & 0 deletions user/apps/tests/dunitest/no_skip.txt
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@ normal/tcp_self_connect_semantics
normal/tcp_dual_stack_semantics
normal/tcp_relisten
normal/tcp_listener_overflow
normal/tcp_accept_handshake
normal/socket_nonblocking
normal/poll_timeout_semantics

normal/epoll_pwait2_semantics
Expand Down
Loading
Loading