diff options
| author | 魏曹先生 <1992414357@qq.com> | 2026-06-20 06:58:31 +0800 |
|---|---|---|
| committer | 魏曹先生 <1992414357@qq.com> | 2026-06-20 06:58:31 +0800 |
| commit | d1cda84f0d9acf2700e4c2eeafff4dbb63cc51fe (patch) | |
| tree | c8c8faff8f419e4d94fae72757d179492ae58b37 /src/output_protocol/ipc.rs | |
| parent | 332ee84b59f33e3a3f1cb109e02a1c0c9fe108cd (diff) | |
feat: implement core capture engine with CLI and output protocols
Diffstat (limited to 'src/output_protocol/ipc.rs')
| -rw-r--r-- | src/output_protocol/ipc.rs | 68 |
1 files changed, 62 insertions, 6 deletions
diff --git a/src/output_protocol/ipc.rs b/src/output_protocol/ipc.rs index c9c88c1..fb6df80 100644 --- a/src/output_protocol/ipc.rs +++ b/src/output_protocol/ipc.rs @@ -1,14 +1,70 @@ use crate::OutputProtocol; +use crate::debug_log; +use std::path::PathBuf; +use std::sync::Arc; +use tokio::io::AsyncWriteExt; +use tokio::net::UnixStream; -#[derive(Debug, Default)] -pub struct IPCOutputProtocol {} +/// IPC output protocol. +/// +/// Connects to a Unix domain socket and sends messages to it. +/// On Windows, this will fail with a clear error message since +/// Unix domain sockets are not supported. +#[derive(Debug)] +pub struct IPCOutputProtocol { + socket_path: PathBuf, + stream: tokio::sync::Mutex<Option<UnixStream>>, +} + +impl IPCOutputProtocol { + pub fn new(socket_path: PathBuf) -> Self { + Self { + socket_path, + stream: tokio::sync::Mutex::new(None), + } + } +} impl OutputProtocol for IPCOutputProtocol { - async fn init(self) { - todo!() + async fn init(&self) { + let path = self.socket_path.clone(); + let stream_lock = &self.stream; + + match UnixStream::connect(&path).await { + Ok(stream) => { + debug_log!("[IPC] Connected to {}", path.display()); + *stream_lock.lock().await = Some(stream); + } + Err(e) => { + debug_log!("[IPC] Failed to connect to {}: {}", path.display(), e); + } + } } - async fn send(self: std::sync::Arc<Self>, _str: &str) { - todo!() + async fn send(self: Arc<Self>, message: &str) { + let mut guard = self.stream.lock().await; + if let Some(ref mut stream) = *guard { + let bytes = format!("{}\n", message); + if stream.write_all(bytes.as_bytes()).await.is_err() { + debug_log!("[IPC] Write error, reconnecting..."); + *guard = None; + // Try to reconnect + if let Ok(new_stream) = UnixStream::connect(&self.socket_path).await { + debug_log!("[IPC] Reconnected to {}", self.socket_path.display()); + *guard = Some(new_stream); + } + } + } else { + // Try to connect + if let Ok(stream) = UnixStream::connect(&self.socket_path).await { + debug_log!("[IPC] Connected to {}", self.socket_path.display()); + let bytes = format!("{}\n", message); + let _ = stream.writable().await; + // We can't use the stream directly here since we need to store it + // for future sends. Let's store and send. + let _ = stream.try_write(bytes.as_bytes()); + *guard = Some(stream); + } + } } } |
