diff --git a/src/cli.rs b/src/cli.rs index 74e7db8..90d1410 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -32,6 +32,18 @@ pub struct ClientCli { /// Path to the scpcap executable on the remote host, e.g. `~/bin/scpcap`. #[arg(long, default_value = "scpcap")] pub remote_exe: String, + + /// Append a diagnostic trace of this (local) process's transfer events + /// to this file -- positions, rename detection, finalize -- for + /// figuring out where a stalled transfer is stuck. + #[arg(long, value_name = "PATH")] + pub log: Option, + + /// Path (on the remote host) to pass as the spawned `--server` + /// process's own `--log`, so both sides of a remote transfer leave a + /// trace. Ignored for local-to-local transfers. + #[arg(long, value_name = "PATH")] + pub remote_log: Option, } /// Remote helper mode, ssh-spawned by the client (like `rsync --server`). @@ -75,6 +87,11 @@ pub struct ServerCli { /// Destination temp path, written while in flight (--recv mode). #[arg(long)] pub dest_temp: Option, + + /// Append a diagnostic trace of this (remote) process's transfer events + /// to this file. See the client's `--log`/`--remote-log`. + #[arg(long, value_name = "PATH")] + pub log: Option, } /// The clap-facing shape of the role choice: two mutually exclusive, diff --git a/src/client.rs b/src/client.rs index 2e4374d..46085fd 100644 --- a/src/client.rs +++ b/src/client.rs @@ -4,9 +4,10 @@ use crate::naming::{self, resolve_dest_paths, strip_name, Target}; use crate::pipeline::{run_recv, run_send, RecvOptions, SendOptions}; use crate::protocol::{MessageSink, MessageSource, WireSink, WireSource}; use crate::ssh::{build_server_command, spawn_remote_server}; +use crate::tracelog::Logger; use anyhow::{bail, Result}; use std::path::PathBuf; -use std::sync::mpsc; +use std::sync::{mpsc, Arc}; use std::thread; pub fn run(cli: ClientCli) -> Result<()> { @@ -16,6 +17,11 @@ pub fn run(cli: ClientCli) -> Result<()> { bail!("unsupported codec {codec:?}: only \"zstd\" is supported"); } + let log: Arc = Arc::new(match &cli.log { + Some(p) => Logger::open(std::path::Path::new(p))?, + None => Logger::none(), + }); + let source_target = Target::parse(&cli.source); let dest_target = Target::parse(&cli.dest); @@ -41,10 +47,12 @@ pub fn run(cli: ClientCli) -> Result<()> { compress: effective_compress, chunk_target, level: cli.level, + log: log.clone(), }; let recv_opts = RecvOptions { dest_final: PathBuf::from(&dest.final_path), dest_temp: PathBuf::from(&dest.temp_path), + log: log.clone(), }; let (tx, rx) = mpsc::channel(); let recv_handle = @@ -54,13 +62,17 @@ pub fn run(cli: ClientCli) -> Result<()> { Ok(()) } (Target::Local(src), Target::Remote { host, .. }) => { - let server_args = vec![ + let mut server_args = vec![ "--recv".to_string(), "--dest-final".to_string(), dest.final_path.clone(), "--dest-temp".to_string(), dest.temp_path.clone(), ]; + if let Some(remote_log) = &cli.remote_log { + server_args.push("--log".to_string()); + server_args.push(remote_log.clone()); + } let remote_cmd = build_server_command(&cli.remote_exe, &server_args); let mut child = spawn_remote_server(host, &remote_cmd)?; let stdin = child.stdin.take().expect("piped stdin"); @@ -72,6 +84,7 @@ pub fn run(cli: ClientCli) -> Result<()> { compress: effective_compress, chunk_target, level: cli.level, + log: log.clone(), }; let sink: Box = Box::new(WireSink(stdin)); let send_result = run_send(send_opts, sink); @@ -101,6 +114,10 @@ pub fn run(cli: ClientCli) -> Result<()> { server_args.push("--compress".to_string()); server_args.push("zstd".to_string()); } + if let Some(remote_log) = &cli.remote_log { + server_args.push("--log".to_string()); + server_args.push(remote_log.clone()); + } let remote_cmd = build_server_command(&cli.remote_exe, &server_args); let mut child = spawn_remote_server(host, &remote_cmd)?; let stdout = child.stdout.take().expect("piped stdout"); @@ -108,6 +125,7 @@ pub fn run(cli: ClientCli) -> Result<()> { let recv_opts = RecvOptions { dest_final: PathBuf::from(&dest.final_path), dest_temp: PathBuf::from(&dest.temp_path), + log: log.clone(), }; let source: Box = Box::new(WireSource(stdout)); let recv_result = run_recv(source, recv_opts); diff --git a/src/main.rs b/src/main.rs index df7e669..00f794b 100644 --- a/src/main.rs +++ b/src/main.rs @@ -8,6 +8,7 @@ mod protocol; mod server; mod ssh; mod tail; +mod tracelog; mod zstd_frame; use clap::Parser; diff --git a/src/pipeline.rs b/src/pipeline.rs index 9557dbb..b71b1ed 100644 --- a/src/pipeline.rs +++ b/src/pipeline.rs @@ -1,12 +1,13 @@ use crate::chunker::{Chunker, RawChunker, RecordAlignedChunker, DEFAULT_FLUSH}; use crate::protocol::{Message, MessageSink, MessageSource}; use crate::tail::{run_tail, SourceEvent}; +use crate::tracelog::Logger; use crate::zstd_frame::{compress_frame, MARKER_FRAME}; use anyhow::{bail, Result}; use std::fs::File; use std::io::Write; use std::path::PathBuf; -use std::sync::mpsc; +use std::sync::{mpsc, Arc}; use std::thread; /// Bound on the reader->compressor and compressor->sender channels: enough to @@ -22,6 +23,7 @@ pub struct SendOptions { pub compress: bool, pub chunk_target: usize, pub level: i32, + pub log: Arc, } enum CompressedItem { @@ -48,7 +50,8 @@ pub fn run_send(opts: SendOptions, mut sink: Box) -> Result<()> let source_path = opts.source_path.clone(); let is_partial = opts.is_partial; let final_source_path = opts.final_source_path.clone(); - move || run_tail(source_path, is_partial, final_source_path, chunker, tail_tx) + let log = opts.log.clone(); + move || run_tail(source_path, is_partial, final_source_path, chunker, tail_tx, log) }); if opts.compress { @@ -112,21 +115,26 @@ pub fn run_send(opts: SendOptions, mut sink: Box) -> Result<()> compressor_result?; } else { let mut send_err: Option = None; + let mut sent: u64 = 0; for event in tail_rx { match event { SourceEvent::Chunk(_kind, data) => { + sent += data.len() as u64; if let Err(e) = sink.send_data(&data) { + opts.log.log(format_args!("send: send_data failed after {sent} bytes: {e}")); send_err = Some(e); break; } } SourceEvent::Eof => { + opts.log.log(format_args!("send: tail EOF, {sent} bytes sent, sending FINALIZE")); if let Err(e) = sink.send_finalize() { send_err = Some(e); } break; } SourceEvent::Error(e) => { + opts.log.log(format_args!("send: tail error after {sent} bytes: {e}")); let _ = sink.send_error(&e); send_err = Some(anyhow::anyhow!(e)); break; @@ -145,6 +153,7 @@ pub fn run_send(opts: SendOptions, mut sink: Box) -> Result<()> pub struct RecvOptions { pub dest_final: PathBuf, pub dest_temp: PathBuf, + pub log: Arc, } /// Writes whatever bytes arrive verbatim to a temp path, then atomically @@ -153,15 +162,33 @@ pub struct RecvOptions { /// error rather than renaming. pub fn run_recv(mut source: Box, opts: RecvOptions) -> Result<()> { let mut f = File::create(&opts.dest_temp)?; + let mut received: u64 = 0; + opts.log.log(format_args!( + "recv: writing to {} (final: {})", + opts.dest_temp.display(), + opts.dest_final.display() + )); loop { match source.recv()? { - Some(Message::Data(payload)) => f.write_all(&payload)?, + Some(Message::Data(payload)) => { + received += payload.len() as u64; + f.write_all(&payload)?; + } Some(Message::Finalize) => { f.sync_all()?; std::fs::rename(&opts.dest_temp, &opts.dest_final)?; + opts.log.log(format_args!( + "recv: FINALIZE after {received} bytes, renamed to {}", + opts.dest_final.display() + )); return Ok(()); } - None => bail!("transfer aborted: stream ended without FINALIZE"), + None => { + opts.log.log(format_args!( + "recv: stream ended without FINALIZE after {received} bytes" + )); + bail!("transfer aborted: stream ended without FINALIZE") + } } } } @@ -190,7 +217,7 @@ mod tests { fn run_pipe(opts: SendOptions, dest_final: PathBuf, dest_temp: PathBuf) -> Result<()> { let (tx, rx) = mpsc::channel::(); - let recv_opts = RecvOptions { dest_final, dest_temp }; + let recv_opts = RecvOptions { dest_final, dest_temp, log: Arc::new(Logger::none()) }; let recv_handle = thread::spawn(move || run_recv(Box::new(ChannelSource(rx)), recv_opts)); run_send(opts, Box::new(ChannelSink(tx)))?; recv_handle.join().unwrap() @@ -211,6 +238,7 @@ mod tests { compress: false, chunk_target: crate::chunker::DEFAULT_TARGET_BYTES, level: crate::zstd_frame::DEFAULT_LEVEL, + log: Arc::new(Logger::none()), }; run_pipe(opts, dest_final.clone(), dest_temp.clone()).unwrap(); @@ -237,13 +265,14 @@ mod tests { compress: false, chunk_target: crate::chunker::DEFAULT_TARGET_BYTES, level: crate::zstd_frame::DEFAULT_LEVEL, + log: Arc::new(Logger::none()), }; let (tx, rx) = mpsc::channel::(); let recv_handle = thread::spawn({ let dest_final = dest_final.clone(); let dest_temp = dest_temp.clone(); - move || run_recv(Box::new(ChannelSource(rx)), RecvOptions { dest_final, dest_temp }) + move || run_recv(Box::new(ChannelSource(rx)), RecvOptions { dest_final, dest_temp, log: Arc::new(Logger::none()) }) }); let send_handle = thread::spawn(move || run_send(opts, Box::new(ChannelSink(tx)))); @@ -275,6 +304,7 @@ mod tests { compress: true, chunk_target: 512, // small target to force multiple frames level: crate::zstd_frame::DEFAULT_LEVEL, + log: Arc::new(Logger::none()), }; run_pipe(opts, dest_final.clone(), dest_temp.clone()).unwrap(); @@ -310,6 +340,7 @@ mod tests { compress: true, chunk_target: crate::chunker::DEFAULT_TARGET_BYTES, level: crate::zstd_frame::DEFAULT_LEVEL, + log: Arc::new(Logger::none()), }; assert!(run_pipe(opts, dest_final.clone(), dest_temp).is_err()); assert!(!dest_final.exists()); diff --git a/src/server.rs b/src/server.rs index 1015a8b..a8b2c42 100644 --- a/src/server.rs +++ b/src/server.rs @@ -2,8 +2,10 @@ use crate::cli::{ServerCli, ServerRole}; use crate::naming; use crate::pipeline::{run_recv, run_send, RecvOptions, SendOptions}; use crate::protocol::{MessageSink, MessageSource, WireSink, WireSource}; +use crate::tracelog::Logger; use anyhow::{bail, Result}; use std::path::PathBuf; +use std::sync::Arc; /// Entry point for `scpcap --server ...`, spawned remotely over ssh by a /// client. Talks the wire protocol over its own stdin/stdout. @@ -14,7 +16,15 @@ pub fn run(cli: ServerCli) -> Result<()> { } } +fn open_log(path: &Option) -> Result> { + Ok(Arc::new(match path { + Some(p) => Logger::open(std::path::Path::new(p))?, + None => Logger::none(), + })) +} + fn run_send_server(cli: ServerCli) -> Result<()> { + let log = open_log(&cli.log)?; let source = cli.source.ok_or_else(|| anyhow::anyhow!("--send requires --source"))?; let final_source = cli .final_source @@ -33,6 +43,7 @@ fn run_send_server(cli: ServerCli) -> Result<()> { compress, chunk_target, level: cli.level, + log, }; // `Stdout`/`Stdin` (not the `.lock()` guards) are used here: the guards // aren't `Send`, but `MessageSink`/`MessageSource` require it. @@ -41,6 +52,7 @@ fn run_send_server(cli: ServerCli) -> Result<()> { } fn run_recv_server(cli: ServerCli) -> Result<()> { + let log = open_log(&cli.log)?; let dest_final = cli .dest_final .ok_or_else(|| anyhow::anyhow!("--recv requires --dest-final"))?; @@ -51,6 +63,7 @@ fn run_recv_server(cli: ServerCli) -> Result<()> { let opts = RecvOptions { dest_final: PathBuf::from(dest_final), dest_temp: PathBuf::from(dest_temp), + log, }; let source: Box = Box::new(WireSource(std::io::stdin())); run_recv(source, opts) diff --git a/src/tail.rs b/src/tail.rs index 8a349ac..d6a4bc6 100644 --- a/src/tail.rs +++ b/src/tail.rs @@ -1,8 +1,10 @@ use crate::chunker::{Chunker, ChunkKind}; +use crate::tracelog::Logger; use anyhow::{bail, Result}; use std::fs::File; use std::os::unix::fs::FileExt; use std::path::{Path, PathBuf}; +use std::sync::Arc; use std::time::{Duration, Instant}; /// Poll interval for growth/rename checks, per STREAMING_ZSTD_FORMAT.md section 9. @@ -50,16 +52,25 @@ pub fn run_tail( final_path: PathBuf, mut chunker: Box, events: impl EventSink, + log: Arc, ) -> Result<()> { if !source_path.exists() { + log.log(format_args!("tail: source not found: {}", source_path.display())); bail!("source not found: {}", source_path.display()); } - let result = drain(&source_path, is_partial, &final_path, chunker.as_mut(), &events); + log.log(format_args!( + "tail: start source={} is_partial={is_partial} final={}", + source_path.display(), + final_path.display() + )); + let result = drain(&source_path, is_partial, &final_path, chunker.as_mut(), &events, &log); match &result { Ok(()) => { + log.log(format_args!("tail: done, sending EOF")); let _ = events.send(SourceEvent::Eof); } Err(e) => { + log.log(format_args!("tail: error: {e}")); let _ = events.send(SourceEvent::Error(e.to_string())); } } @@ -72,10 +83,13 @@ fn drain( final_path: &Path, chunker: &mut dyn Chunker, events: &impl EventSink, + log: &Logger, ) -> Result<()> { let file = File::open(source_path)?; let mut pos: u64 = 0; let mut complete_signal = !is_partial; + let mut shrunk_below_pos = false; + let mut last_wait_log: Option = None; let mut buf = vec![0u8; READ_WINDOW]; loop { @@ -94,6 +108,7 @@ fn drain( pos += n as u64; emit_ready(chunker, events)?; made_progress = true; + shrunk_below_pos = false; } } else if chunker.should_flush(Instant::now()) && let Some((kind, data)) = chunker.flush() @@ -103,11 +118,22 @@ fn drain( // If file_len < pos, the writer truncated/restarted mid-file // (STREAMING_ZSTD_FORMAT.md section 7.7): do nothing, never rewind // `pos`, just wait for it to regrow past `pos` on a later poll. + if file_len < pos && !shrunk_below_pos { + shrunk_below_pos = true; + log.log(format_args!( + "tail: file shrank below pos (pos={pos}, file_len={file_len}); waiting for regrowth -- \ + if the file never grows past pos again, this transfer will never finish" + )); + } if !complete_signal && final_path.try_exists()? { // Rename observed -- signal now, well before finalization, per // section 7.4's repoint-then-drain-then-finalize ordering. complete_signal = true; + log.log(format_args!( + "tail: rename to {} observed at pos={pos}, file_len={file_len}", + final_path.display() + )); } let current_len = file.metadata()?.len(); @@ -115,8 +141,19 @@ fn drain( if let Some((kind, data)) = chunker.finish()? { send_chunk(events, kind, data)?; } + log.log(format_args!("tail: finalized at pos={pos}")); return Ok(()); } + if complete_signal + && pos != current_len + && last_wait_log.is_none_or(|t| t.elapsed() >= Duration::from_secs(1)) + { + last_wait_log = Some(Instant::now()); + log.log(format_args!( + "tail: waiting to finalize: pos={pos} != current_len={current_len}{}", + if current_len < pos { " (final file is smaller than what was already read -- this will never catch up)" } else { "" } + )); + } if !made_progress { std::thread::sleep(POLL_INTERVAL); @@ -165,7 +202,7 @@ mod tests { std::fs::write(&path, b"hello world").unwrap(); let (tx, rx) = mpsc::channel(); - run_tail(path.clone(), false, path.clone(), Box::new(RawChunker::new()), tx).unwrap(); + run_tail(path.clone(), false, path.clone(), Box::new(RawChunker::new()), tx, Arc::new(Logger::none())).unwrap(); let (data, eof) = collect(rx); assert_eq!(data, b"hello world"); assert!(eof); @@ -177,7 +214,7 @@ mod tests { let path = dir.path().join("nope.pcap.partial"); let final_path = dir.path().join("nope.pcap"); let (tx, _rx) = mpsc::channel(); - let err = run_tail(path, true, final_path, Box::new(RawChunker::new()), tx).unwrap_err(); + let err = run_tail(path, true, final_path, Box::new(RawChunker::new()), tx, Arc::new(Logger::none())).unwrap_err(); assert!(err.to_string().contains("not found")); } @@ -194,7 +231,7 @@ mod tests { let handle = std::thread::spawn({ let partial_path = partial_path.clone(); let final_path = final_path.clone(); - move || run_tail(partial_path, true, final_path, Box::new(RawChunker::new()), tx) + move || run_tail(partial_path, true, final_path, Box::new(RawChunker::new()), tx, Arc::new(Logger::none())) }); std::thread::sleep(Duration::from_millis(60)); @@ -210,6 +247,34 @@ mod tests { assert!(eof); } + #[test] + fn log_records_rename_detection_and_finalize() { + let dir = tempdir().unwrap(); + let partial_path = dir.path().join("cap.pcap.partial"); + let final_path = dir.path().join("cap.pcap"); + std::fs::write(&partial_path, b"chunk1").unwrap(); + let log_path = dir.path().join("trace.log"); + let log = Arc::new(Logger::open(&log_path).unwrap()); + + let (tx, rx) = mpsc::channel(); + let handle = std::thread::spawn({ + let partial_path = partial_path.clone(); + let final_path = final_path.clone(); + let log = log.clone(); + move || run_tail(partial_path, true, final_path, Box::new(RawChunker::new()), tx, log) + }); + + std::thread::sleep(Duration::from_millis(60)); + std::fs::rename(&partial_path, &final_path).unwrap(); + + handle.join().unwrap().unwrap(); + collect(rx); + + let contents = std::fs::read_to_string(&log_path).unwrap(); + assert!(contents.contains("rename to"), "log missing rename line: {contents}"); + assert!(contents.contains("finalized"), "log missing finalize line: {contents}"); + } + #[test] fn held_fd_survives_rename_and_even_unlink() { let dir = tempdir().unwrap(); @@ -223,7 +288,7 @@ mod tests { let handle = std::thread::spawn({ let partial_path = partial_path.clone(); let final_path = final_path.clone(); - move || run_tail(partial_path, true, final_path, Box::new(RawChunker::new()), tx) + move || run_tail(partial_path, true, final_path, Box::new(RawChunker::new()), tx, Arc::new(Logger::none())) }); std::thread::sleep(Duration::from_millis(60)); @@ -254,7 +319,7 @@ mod tests { let handle = std::thread::spawn({ let partial_path = partial_path.clone(); let final_path = final_path.clone(); - move || run_tail(partial_path, true, final_path, Box::new(RawChunker::new()), tx) + move || run_tail(partial_path, true, final_path, Box::new(RawChunker::new()), tx, Arc::new(Logger::none())) }); // Let the reader fully drain the initial 10 bytes. diff --git a/src/tracelog.rs b/src/tracelog.rs new file mode 100644 index 0000000..799f205 --- /dev/null +++ b/src/tracelog.rs @@ -0,0 +1,30 @@ +use std::fs::{File, OpenOptions}; +use std::io::Write; +use std::path::Path; +use std::sync::Mutex; +use std::time::{SystemTime, UNIX_EPOCH}; + +/// Optional diagnostic trace, one timestamped line per event, for diagnosing +/// where a transfer stalls (e.g. which side is still waiting, and on what). +/// `Logger::none()` is a no-op so call sites don't need to branch on whether +/// `--log` was given. +pub struct Logger(Option>); + +impl Logger { + pub fn none() -> Logger { + Logger(None) + } + + pub fn open(path: &Path) -> std::io::Result { + let f = OpenOptions::new().create(true).append(true).open(path)?; + Ok(Logger(Some(Mutex::new(f)))) + } + + pub fn log(&self, args: std::fmt::Arguments) { + let Some(m) = &self.0 else { return }; + let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default(); + if let Ok(mut f) = m.lock() { + let _ = writeln!(f, "[{:.6}] {}", now.as_secs_f64(), args); + } + } +}