fix(debug-ipc): 修复连接生命周期和完整帧写入
区分 IPC server 所有者与普通克隆,避免临时 clone 析构时提前停止共享服务。 改用带超时的阻塞 write_all 发送完整数据帧,防止非阻塞短写破坏热重载协议。
This commit is contained in:
@@ -11,9 +11,18 @@ use std::{
|
|||||||
|
|
||||||
use crate::error::{CliError, Result};
|
use crate::error::{CliError, Result};
|
||||||
|
|
||||||
#[derive(Clone)]
|
|
||||||
pub struct DebugIpcServer {
|
pub struct DebugIpcServer {
|
||||||
inner: Arc<Inner>,
|
inner: Arc<Inner>,
|
||||||
|
owns_lifecycle: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Clone for DebugIpcServer {
|
||||||
|
fn clone(&self) -> Self {
|
||||||
|
Self {
|
||||||
|
inner: Arc::clone(&self.inner),
|
||||||
|
owns_lifecycle: false,
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
struct Inner {
|
struct Inner {
|
||||||
@@ -39,7 +48,10 @@ impl DebugIpcServer {
|
|||||||
let handle = thread::spawn(move || accept_loop(listener, thread_inner));
|
let handle = thread::spawn(move || accept_loop(listener, thread_inner));
|
||||||
*inner.thread.lock().unwrap() = Some(handle);
|
*inner.thread.lock().unwrap() = Some(handle);
|
||||||
|
|
||||||
Ok(Self { inner })
|
Ok(Self {
|
||||||
|
inner,
|
||||||
|
owns_lifecycle: true,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn port(&self) -> u16 {
|
pub fn port(&self) -> u16 {
|
||||||
@@ -50,12 +62,11 @@ impl DebugIpcServer {
|
|||||||
let frame = encode_frame(message_type, payload)?;
|
let frame = encode_frame(message_type, payload)?;
|
||||||
let mut clients = self.inner.clients.lock().unwrap();
|
let mut clients = self.inner.clients.lock().unwrap();
|
||||||
let mut any = false;
|
let mut any = false;
|
||||||
clients.retain_mut(|stream| match stream.write(&frame) {
|
clients.retain_mut(|stream| match stream.write_all(&frame) {
|
||||||
Ok(_) => {
|
Ok(()) => {
|
||||||
any = true;
|
any = true;
|
||||||
true
|
true
|
||||||
}
|
}
|
||||||
Err(err) if err.kind() == io::ErrorKind::WouldBlock => true,
|
|
||||||
Err(_) => false,
|
Err(_) => false,
|
||||||
});
|
});
|
||||||
Ok(any)
|
Ok(any)
|
||||||
@@ -81,8 +92,10 @@ impl DebugIpcServer {
|
|||||||
|
|
||||||
impl Drop for DebugIpcServer {
|
impl Drop for DebugIpcServer {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
|
if self.owns_lifecycle {
|
||||||
self.stop();
|
self.stop();
|
||||||
}
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn encode_frame(message_type: u16, payload: &[u8]) -> Result<Vec<u8>> {
|
pub fn encode_frame(message_type: u16, payload: &[u8]) -> Result<Vec<u8>> {
|
||||||
@@ -99,9 +112,13 @@ fn accept_loop(listener: TcpListener, inner: Arc<Inner>) {
|
|||||||
while !inner.stop.load(Ordering::Relaxed) {
|
while !inner.stop.load(Ordering::Relaxed) {
|
||||||
match listener.accept() {
|
match listener.accept() {
|
||||||
Ok((stream, _)) => {
|
Ok((stream, _)) => {
|
||||||
let _ = stream.set_nonblocking(true);
|
if stream
|
||||||
|
.set_write_timeout(Some(Duration::from_secs(1)))
|
||||||
|
.is_ok()
|
||||||
|
{
|
||||||
inner.clients.lock().unwrap().push(stream);
|
inner.clients.lock().unwrap().push(stream);
|
||||||
}
|
}
|
||||||
|
}
|
||||||
Err(err) if err.kind() == io::ErrorKind::WouldBlock => {
|
Err(err) if err.kind() == io::ErrorKind::WouldBlock => {
|
||||||
thread::sleep(Duration::from_millis(50));
|
thread::sleep(Duration::from_millis(50));
|
||||||
}
|
}
|
||||||
@@ -112,7 +129,7 @@ fn accept_loop(listener: TcpListener, inner: Arc<Inner>) {
|
|||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::encode_frame;
|
use super::{DebugIpcServer, encode_frame};
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn encodes_big_endian_ipc_frame() {
|
fn encodes_big_endian_ipc_frame() {
|
||||||
@@ -120,4 +137,15 @@ mod tests {
|
|||||||
assert_eq!(&frame[..6], &[0, 2, 0, 0, 0, 7]);
|
assert_eq!(&frame[..6], &[0, 2, 0, 0, 0, 7]);
|
||||||
assert_eq!(&frame[6..], b"[\"a.b\"]");
|
assert_eq!(&frame[6..], b"[\"a.b\"]");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn dropping_clone_keeps_server_running() {
|
||||||
|
let server = DebugIpcServer::start().unwrap();
|
||||||
|
let port = server.port();
|
||||||
|
|
||||||
|
drop(server.clone());
|
||||||
|
|
||||||
|
assert_eq!(server.port(), port);
|
||||||
|
server.safe_exit();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user