diff --git a/src/commands.rs b/src/commands.rs index 58c04cccd..805827dd3 100644 --- a/src/commands.rs +++ b/src/commands.rs @@ -15,13 +15,13 @@ use crate::cache::{IpcStorage, storage_from_config}; use crate::client::{ServerConnection, connect_to_server, connect_with_retry}; use crate::cmdline::{Command, StatsFormat}; -use crate::compiler::ColorMode; +use crate::compiler::{ColorMode, is_rmeta_artifact_notification}; use crate::config::{Config, default_disk_cache_dir}; use crate::jobserver::Client; use crate::mock_command::{CommandChild, CommandCreatorSync, ProcessCommandCreator, RunCommand}; use crate::protocol::{Compile, CompileFinished, CompileResponse, Request, Response}; use crate::server::{self, DistInfo, ServerInfo, ServerStartup, ServerStats}; -use crate::util::{daemonize, new_client_runtime}; +use crate::util::{daemonize, new_client_runtime, run_input_output}; use byteorder::{BigEndian, ByteOrder}; use fs::{File, OpenOptions}; use fs_err as fs; @@ -528,35 +528,42 @@ fn handle_compile_response( where T: CommandCreatorSync, { + let mut forwarded_notification = false; let result = match &response { CompileResponse::CompileStarted => { debug!("Server sent CompileStarted"); - // Wait for CompileFinished. - match conn.read_one_response() { - Ok(Response::CompileFinished(result)) => Some(result), - Ok(_) => bail!("unexpected response from server"), - Err(e) => { - match e.downcast_ref::() { - Some(io_e) if io_e.kind() == io::ErrorKind::UnexpectedEof => { - eprintln!( - "sccache: warning: The server looks like it shut down \ - unexpectedly, compiling locally instead" - ); - } - _ => { - //TODO: something better here? - if ignore_all_server_io_errors() { + loop { + match conn.read_one_response() { + Ok(Response::ArtifactNotification(chunk)) => { + stderr.write_all(&chunk)?; + stderr.flush()?; + forwarded_notification = true; + } + Ok(Response::CompileFinished(result)) => break Some(result), + Ok(_) => bail!("unexpected response from server"), + Err(e) => { + match e.downcast_ref::() { + Some(io_e) if io_e.kind() == io::ErrorKind::UnexpectedEof => { eprintln!( - "sccache: warning: error reading compile response from server \ - compiling locally instead" + "sccache: warning: The server looks like it shut down \ + unexpectedly, compiling locally instead" ); - } else { - return Err(e) - .context("error reading compile response from server"); + } + _ => { + //TODO: something better here? + if ignore_all_server_io_errors() { + eprintln!( + "sccache: warning: error reading compile response from server \ + compiling locally instead" + ); + } else { + return Err(e) + .context("error reading compile response from server"); + } } } + break None; } - None } } } @@ -564,18 +571,32 @@ where }; handle_compile_result( - creator, runtime, response, result, exe, cmdline, cwd, stdout, stderr, + creator, + runtime, + response, + result, + forwarded_notification, + exe, + cmdline, + cwd, + stdout, + stderr, ) } /// Dispatch the outcome of a compile, whether received from the daemon over IPC /// or produced by a local `SccacheService` in client-side mode. +/// +/// * `forwarded_notification` +/// is whether a [`Response::ArtifactNotification`] already reached `stderr`, +/// in which case its copy in compiler's stderr is dropped. #[allow(clippy::too_many_arguments)] fn handle_compile_result( mut creator: T, runtime: &mut Runtime, response: CompileResponse, finished: Option, + forwarded_notification: bool, exe: &Path, cmdline: Vec, cwd: &Path, @@ -587,7 +608,10 @@ where { match response { CompileResponse::CompileStarted => { - if let Some(finished) = finished { + if let Some(mut finished) = finished { + if forwarded_notification { + finished.stderr = without_rmeta_artifact_notifications(&finished.stderr); + } return handle_compile_finished(finished, stdout, stderr); } // Server disconnected before sending CompileFinished; fall back to local compilation. @@ -607,13 +631,30 @@ where trace!("running command: {:?}", cmd); } - let status = runtime.block_on(async move { - let child = cmd.spawn().await?; - child - .wait() - .await - .with_context(|| "failed to wait for a child") - })?; + let status = if !forwarded_notification { + runtime.block_on(async move { + let child = cmd.spawn().await?; + child + .wait() + .await + .with_context(|| "failed to wait for a child") + })? + } else { + // The `.rmeta` notification was already forwarded so drop them from local rustc's copy. + // Assumption: this rewrites artifacts which callers may be reading, + // which is safe only while rustc output is deterministic (it is as of 1.98.1) + let output = match runtime.block_on(run_input_output(cmd, None)) { + Ok(output) => output, + Err(e) => match e.downcast::() { + Ok(ProcessError(output)) => output, + Err(e) => return Err(e).context("failed to wait for a child"), + }, + }; + stdout.write_all(&output.stdout)?; + stderr.write_all(&without_rmeta_artifact_notifications(&output.stderr))?; + stderr.flush()?; + output.status + }; Ok(status.code().unwrap_or_else(|| { if let Some(sig) = status_signal(status) { @@ -624,6 +665,16 @@ where })) } +/// Drop `.rmeta` artifact notification lines from `stderr` +/// when they have already been forwarded to the caller. +fn without_rmeta_artifact_notifications(stderr: &[u8]) -> Vec { + stderr + .split_inclusive(|&b| b == b'\n') + .filter(|line| !is_rmeta_artifact_notification(line)) + .collect::>() + .concat() +} + /// Send a `Compile` request to the sccache server `conn`, and handle the response. /// /// The first entry in `cmdline` will be looked up in `path` if it is not @@ -694,13 +745,15 @@ where args: cmdline.clone(), env_vars, }; - let (compile_resp, finished) = runtime.block_on(service.compile_direct(compile))?; + let (compile_resp, finished, forwarded_notification) = + runtime.block_on(service.compile_direct(compile, &mut *stderr))?; let creator = C::new(jobserver); let exit_code = handle_compile_result( creator, runtime, compile_resp, finished, + forwarded_notification, &exe_path, cmdline, cwd, diff --git a/src/compiler/compiler.rs b/src/compiler/compiler.rs index 56dea23c0..ea22de6b7 100644 --- a/src/compiler/compiler.rs +++ b/src/compiler/compiler.rs @@ -36,7 +36,8 @@ use crate::lru_disk_cache; use crate::mock_command::{CommandChild, CommandCreatorSync, RunCommand, exit_status}; use crate::server; use crate::util::{ - Digest, fmt_duration_as_secs, resolve_compiler_avoiding_wrapper, run_input_output, + Digest, StderrLineObserver, fmt_duration_as_secs, resolve_compiler_avoiding_wrapper, + run_input_output, run_input_output_observing, }; use crate::{counted_array, dist}; use async_trait::async_trait; @@ -196,7 +197,7 @@ impl CompileCommandImpl for SingleCompileCommand { async fn execute( &self, - _: &server::SccacheService, + service: &server::SccacheService, creator: &T, ) -> Result where @@ -219,7 +220,18 @@ impl CompileCommandImpl for SingleCompileCommand { if *share_jobserver { cmd.share_jobserver(); } - run_input_output(cmd, None).await + match service.notification_tx() { + Some(tx) => { + let tx = tx.clone(); + let observer: StderrLineObserver = Box::new(move |line| { + if super::rust::is_rmeta_artifact_notification(line) { + let _ = tx.unbounded_send(line.to_vec()); + } + }); + run_input_output_observing(cmd, None, Some(observer)).await + } + None => run_input_output(cmd, None).await, + } } } diff --git a/src/compiler/mod.rs b/src/compiler/mod.rs index 46a16a46d..926c15882 100644 --- a/src/compiler/mod.rs +++ b/src/compiler/mod.rs @@ -36,3 +36,4 @@ mod counted_array; pub use crate::compiler::c::CCompilerKind; pub use crate::compiler::compiler::*; pub use crate::compiler::preprocessor_cache::PreprocessorCacheEntry; +pub use crate::compiler::rust::is_rmeta_artifact_notification; diff --git a/src/compiler/rust.rs b/src/compiler/rust.rs index 21b7e9450..9c5dd6d82 100644 --- a/src/compiler/rust.rs +++ b/src/compiler/rust.rs @@ -2752,6 +2752,23 @@ fn parse_rustc_z_ls(stdout: &str) -> Result> { Ok(dep_names) } +/// Whether `line` is rustc's `--json=artifacts` notification for a `.rmeta`. +pub fn is_rmeta_artifact_notification(line: &[u8]) -> bool { + // Cheap reject before parsing as most stderr lines are diagnostics. + if memchr::memmem::find(line, b".rmeta").is_none() { + return false; + } + #[derive(serde::Deserialize)] + struct ArtifactNotification<'a> { + // `Cow` as Windows path may need escapes + #[serde(borrow)] + artifact: std::borrow::Cow<'a, str>, + } + serde_json::from_slice::>(line) + .map(|n| n.artifact.ends_with(".rmeta")) + .unwrap_or(false) +} + #[cfg(test)] mod test { use super::*; @@ -4139,4 +4156,35 @@ proc_macro false let result = parse_arguments(&args, cwd); assert!(matches!(result, CompilerArguments::CannotCache(..))); } + + #[test] + fn test_is_rmeta_artifact_notification() { + // rmeta notification + assert!(is_rmeta_artifact_notification( + br#"{"artifact":"/t/deps/libdep-1234.rmeta","emit":"metadata"}"# + )); + assert!(is_rmeta_artifact_notification( + b"{\"artifact\":\"/t/deps/libdep-1234.rmeta\",\"emit\":\"metadata\"}\n" + )); + // Paths with JSON escapes: Windows separators, or an escaped quote. + assert!(is_rmeta_artifact_notification( + br#"{"artifact":"C:\\t\\deps\\libdep-1234.rmeta","emit":"metadata"}"# + )); + assert!(is_rmeta_artifact_notification( + br#"{"artifact":"/t/de\"ps/libdep-1234.rmeta","emit":"metadata"}"# + )); + + // not rmeta notification + assert!(!is_rmeta_artifact_notification( + br#"{"artifact":"/t/deps/libdep-1234.rlib","emit":"link"}"# + )); + assert!(!is_rmeta_artifact_notification( + br#"{"artifact":"/t/deps/dep-1234.d","emit":"dep-info"}"# + )); + assert!(!is_rmeta_artifact_notification( + br#"{"$message_type":"diagnostic","message":"unused variable","level":"warning"}"# + )); + assert!(!is_rmeta_artifact_notification(b"not json\n")); + assert!(!is_rmeta_artifact_notification(b"")); + } } diff --git a/src/protocol.rs b/src/protocol.rs index 7b3d8205f..a35c1388b 100644 --- a/src/protocol.rs +++ b/src/protocol.rs @@ -68,6 +68,14 @@ pub enum Response { StoragePutPreprocessorEntry(Result<(), String>), /// Response for `Request::RecordStats`. RecordStats, + + /// A rustc artifact notification delivered while the compiler is still running. + /// + /// The line also stays in [`CompileFinished::stderr`]. + /// The client that has forwarded the notification must drop that copy. + /// + /// Currently only `.rmeta` notifications are sent (for Cargo pipelining). + ArtifactNotification(Vec), } /// Possible responses from the server for a `Compile` request. diff --git a/src/server.rs b/src/server.rs index 2f142684a..e216d9d28 100644 --- a/src/server.rs +++ b/src/server.rs @@ -826,10 +826,18 @@ where /// This field causes [WaitUntilZero] to wait until this struct drops. #[allow(dead_code)] info: ActiveInfo, + + /// Notifications to deliver to the client while compiler is still running, + /// as [`Response::ArtifactNotification`], before [`CompileFinished`]. + /// + /// * `Some` only on the clone [`SccacheService::check_compiler`] gives a + /// rustc compile task. + /// * `None` on the shared service and every other clone. + notification_tx: Option>>, } type SccacheRequest = Message>; -type SccacheResponse = Message> + Send>>>; +type SccacheResponse = Message> + Send>>>; /// Messages sent from all services to the main event loop indicating activity. /// @@ -1012,6 +1020,7 @@ where creator: C::new(client), tx, info, + notification_tx: None, } } @@ -1033,6 +1042,7 @@ where creator: C::new(&client), tx, info, + notification_tx: None, } } @@ -1068,6 +1078,7 @@ where creator: C::new(&client), tx, info, + notification_tx: None, } } @@ -1109,12 +1120,14 @@ where Message::WithoutBody(message) => { sink.send(Frame::Message { message }).await?; } - Message::WithBody(message, body) => { + Message::WithBody(message, mut body) => { sink.send(Frame::Message { message }).await?; - sink.send(Frame::Body { - chunk: Some(util::spawn(body).await??), - }) - .await?; + while let Some(chunk) = body.next().await { + sink.send(Frame::Body { + chunk: Some(chunk?), + }) + .await?; + } sink.send(Frame::Body { chunk: None }).await?; } } @@ -1155,6 +1168,14 @@ where std::mem::take(&mut *s) } + /// Artifact notifications to deliver to the client while the compiler is + /// still running before [`CompileFinished`]. + /// + /// See [`SccacheService::notification_tx`] for more. + pub fn notification_tx(&self) -> Option<&futures::channel::mpsc::UnboundedSender>> { + self.notification_tx.as_ref() + } + async fn merge_stats(&self, delta: ServerStats) { *self.stats.lock().await += delta; } @@ -1179,21 +1200,36 @@ where /// Run a compile entirely in the current process (used in client-side mode). /// - /// Returns the `CompileResponse` variant and, when compilation started, the - /// accompanying `CompileFinished` result. + /// Returns a tuple of + /// + /// * the `CompileResponse` variant, and + /// * when compilation started, the accompanying `CompileFinished` result, and + /// * whether a notification was written to `stderr` while the compiler was running. pub async fn compile_direct( &self, compile: Compile, - ) -> Result<(CompileResponse, Option)> { + stderr: &mut dyn Write, + ) -> Result<(CompileResponse, Option, bool)> { match self.handle_compile(compile).await? { - Message::WithBody(Response::Compile(resp), body) => { - let finished = match body.await? { - Response::CompileFinished(f) => f, - _ => bail!("unexpected body response from compile_direct"), - }; - Ok((resp, Some(finished))) + Message::WithBody(Response::Compile(resp), mut body) => { + let mut finished = None; + let mut forwarded_notification = false; + while let Some(item) = body.next().await { + match item? { + Response::ArtifactNotification(chunk) => { + stderr.write_all(&chunk)?; + stderr.flush()?; + forwarded_notification = true; + } + Response::CompileFinished(f) => finished = Some(f), + _ => bail!("unexpected body response from compile_direct"), + } + } + let finished = finished + .ok_or_else(|| anyhow!("compile body ended without CompileFinished"))?; + Ok((resp, Some(finished), forwarded_notification)) } - Message::WithoutBody(Response::Compile(resp)) => Ok((resp, None)), + Message::WithoutBody(Response::Compile(resp)) => Ok((resp, None, false)), _ => bail!("unexpected response from handle_compile in compile_direct"), } } @@ -1380,11 +1416,29 @@ where CompilerArguments::Ok(hasher) => { debug!("parse_arguments: Ok: {:?}", cmd); - let body = self - .clone() - .start_compile_task(c, hasher, cmd, cwd, env_vars) - .and_then(|res| async { Ok(Response::CompileFinished(res)) }) - .boxed(); + let body = if c.kind() == CompilerKind::Rust { + // Only rustc emits artifact notifications, + // for `.rmeta` pipelining. + let (tx, rx) = futures::channel::mpsc::unbounded(); + let mut me = self.clone(); + me.notification_tx = Some(tx); + let finished = util::spawn_on( + &self.rt, + me.start_compile_task(c, hasher, cmd, cwd, env_vars), + ); + rx.map(|chunk| Ok(Response::ArtifactNotification(chunk))) + .chain(futures::stream::once(async move { + finished.await?.map(Response::CompileFinished) + })) + .boxed() + } else { + futures::stream::once( + self.clone() + .start_compile_task(c, hasher, cmd, cwd, env_vars) + .map_ok(Response::CompileFinished), + ) + .boxed() + }; return Message::WithBody( Response::Compile(CompileResponse::CompileStarted), @@ -2456,6 +2510,8 @@ fn waits_until_zero() { #[cfg(test)] mod tests { use super::*; + use crate::mock_command::{MockChild, MockCommandCreator, exit_status}; + use crate::test::utils::{TestFixture, next_command, next_command_calls}; struct StringWriter { buffer: String, @@ -2584,4 +2640,182 @@ mod tests { assert!(find_s1 < find_s2); } } + + const RMETA_NOTIFICATION: &[u8] = b"{\"artifact\":\"/t/libdep.rmeta\",\"emit\":\"metadata\"}\n"; + const OTHER_STDERR: &[u8] = b"{\"artifact\":\"/t/libdep.rlib\",\"emit\":\"link\"}\n"; + + type MockCreator = Arc>; + + fn mock_rustc_cache_miss(creator: &MockCreator, f: &TestFixture) { + // rustc -vV + next_command( + creator, + Ok(MockChild::new( + exit_status(0), + "rustc 1.90.0 (0000000 2026-01-01)\nhost: x86_64-unknown-linux-gnu\n", + "", + )), + ); + // rustc +stable: not a rustup proxy + next_command(creator, Ok(MockChild::new(exit_status(1), "", ""))); + // rustc --print=sysroot + next_command( + creator, + Ok(MockChild::new( + exit_status(0), + f.tempdir.path().to_str().unwrap(), + "", + )), + ); + mock_rustc_hash_inputs(creator); + // The compile itself. + let out_dir = f.tempdir.path().to_path_buf(); + next_command_calls(creator, move |_| { + std::fs::write(out_dir.join("libdep.rlib"), "rlib")?; + std::fs::write(out_dir.join("libdep.rmeta"), "rmeta")?; + Ok(MockChild::new( + exit_status(0), + "", + [RMETA_NOTIFICATION, OTHER_STDERR].concat(), + )) + }); + } + + fn mock_rustc_hash_inputs(creator: &MockCreator) { + // rustc --emit dep-info -o + next_command_calls(creator, |args| { + let dep_file = args.iter().skip_while(|a| *a != "-o").nth(1).unwrap(); + std::fs::write(dep_file, "libdep.rlib: lib.rs\nlib.rs:\n")?; + Ok(MockChild::new(exit_status(0), "", "")) + }); + // rustc --print file-names + next_command( + creator, + Ok(MockChild::new(exit_status(0), "libdep.rlib\n", "")), + ); + } + + type MockService = SccacheService; + + /// A service whose mocked rustc is queued for one cache-miss compile of `dep`, + /// and then a request for that compile. + fn rustc_miss_fixture() -> (TestFixture, Runtime, MockService, Compile) { + use crate::cache::disk::DiskCache; + use crate::config::PreprocessorCacheModeConfig; + + let _ = env_logger::try_init(); + let f = TestFixture::new(); + let rustc = f.mk_bin("rustc").unwrap(); + // Windows uses bin, everything else uses lib. Just create both. + std::fs::create_dir(f.tempdir.path().join("lib")).unwrap(); + std::fs::create_dir(f.tempdir.path().join("bin")).unwrap(); + f.touch("lib.rs").unwrap(); + + let runtime = Runtime::new().unwrap(); + let storage = Arc::new(DiskCache::new( + f.tempdir.path().join("cache"), + u64::MAX, + runtime.handle(), + PreprocessorCacheModeConfig::default(), + CacheMode::ReadWrite, + vec![], + )); + let service = MockService::mock_with_storage(storage, runtime.handle().clone()); + mock_rustc_cache_miss(&service.creator, &f); + + let compile = dep_compile(&f, &rustc); + (f, runtime, service, compile) + } + + /// The request for compiling `dep`. + fn dep_compile(f: &TestFixture, rustc: &std::path::Path) -> Compile { + Compile { + exe: rustc.into(), + cwd: f.tempdir.path().into(), + args: [ + "--crate-name", + "dep", + "lib.rs", + "--crate-type", + "lib", + "--emit=link,metadata", + "--out-dir", + f.tempdir.path().to_str().unwrap(), + ] + .into_iter() + .map(Into::into) + .collect(), + env_vars: vec![], + } + } + + fn compile_body( + runtime: &Runtime, + service: &MockService, + compile: Compile, + ) -> (Vec, CompileFinished) { + let (resp, body) = match runtime.block_on(service.handle_compile(compile)).unwrap() { + Message::WithBody(Response::Compile(resp), body) => (resp, body), + other => panic!("unexpected response: {:?}", other.into_inner()), + }; + assert!(matches!(resp, CompileResponse::CompileStarted)); + + let mut body = runtime.block_on(body.try_collect::>()).unwrap(); + let finished = match body.pop() { + Some(Response::CompileFinished(finished)) => finished, + other => panic!("body must end with CompileFinished, got {other:?}"), + }; + (body, finished) + } + + #[test] + fn daemon_mode_rmeta_notification_delivery_on_rustc_miss() { + let (_f, runtime, service, compile) = rustc_miss_fixture(); + + let (before_finished, finished) = compile_body(&runtime, &service, compile); + + assert_eq!(Some(0), finished.retcode); + assert!( + matches!( + before_finished.as_slice(), + [Response::ArtifactNotification(chunk)] if chunk == RMETA_NOTIFICATION + ), + "delivered while the compiler runs: {before_finished:?}" + ); + // The daemon leaves rustc's stderr intact. + // Client will dedup stderr if already forwarded. + assert_eq!( + [RMETA_NOTIFICATION, OTHER_STDERR].concat(), + finished.stderr, + "stderr delivered with CompileFinished" + ); + assert_eq!(0, service.creator.lock().unwrap().children.len()); + } + + #[test] + fn client_side_mode_rmeta_notification_delivery_on_rustc_miss() { + let (_f, runtime, service, compile) = rustc_miss_fixture(); + + let mut stderr = Vec::new(); + let (resp, finished, forwarded) = runtime + .block_on(service.compile_direct(compile, &mut stderr)) + .unwrap(); + let finished = finished.expect("compile started"); + + assert!(matches!(resp, CompileResponse::CompileStarted)); + assert_eq!(Some(0), finished.retcode); + assert!(forwarded); + assert_eq!( + RMETA_NOTIFICATION, + stderr.as_slice(), + "forwarded while the compiler runs" + ); + // The caller drops the copy in `CompileFinished`. + assert_eq!( + [RMETA_NOTIFICATION, OTHER_STDERR].concat(), + finished.stderr, + "stderr delivered with CompileFinished" + ); + assert_eq!(0, service.creator.lock().unwrap().children.len()); + } } diff --git a/src/test/tests.rs b/src/test/tests.rs index ca63aa92b..b9ed9afcd 100644 --- a/src/test/tests.rs +++ b/src/test/tests.rs @@ -294,3 +294,201 @@ fn test_server_compile() { // Ensure that it shuts down. child.join().unwrap(); } + +/// A [`Write`] sink shared with the fake server. +/// +/// This is needed so it can observe how much the client has written +/// at a given point in the protocol exchange. +#[derive(Clone, Default)] +struct SharedWriter(Arc>>); + +impl Write for SharedWriter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().unwrap().extend_from_slice(buf); + Ok(buf.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + +const RMETA_NOTIFICATION: &[u8] = b"{\"artifact\":\"/t/libdep.rmeta\",\"emit\":\"metadata\"}\n"; +const OTHER_STDERR: &[u8] = b"{\"artifact\":\"/t/libdep.rlib\",\"emit\":\"link\"}\n"; + +/// Accepts one connection as a fake daemon. +/// +/// The fake one replies to a compile request until a `.rmeta` notification comes. +/// It leaves it for caller to continue working on it. +fn fake_daemon_until_rmeta_notification(listener: std::net::TcpListener) -> std::net::TcpStream { + use crate::protocol::{CompileResponse, Request, Response}; + use crate::util::write_length_prefixed_bincode; + use byteorder::{BigEndian, ByteOrder}; + use std::io::Read; + + let (mut sock, _) = listener.accept().unwrap(); + let mut len = [0; 4]; + sock.read_exact(&mut len).unwrap(); + let mut req = vec![0; BigEndian::read_u32(&len) as usize]; + sock.read_exact(&mut req).unwrap(); + assert!(matches!( + bincode::deserialize::(&req).unwrap(), + Request::Compile(_) + )); + + write_length_prefixed_bincode( + &mut sock, + Response::Compile(CompileResponse::CompileStarted), + ) + .unwrap(); + write_length_prefixed_bincode( + &mut sock, + Response::ArtifactNotification(RMETA_NOTIFICATION.to_vec()), + ) + .unwrap(); + sock +} + +/// Runs a client compiling `dep` against daemon at `addr`. +fn compile_dep( + creator: Arc>, + f: &TestFixture, + addr: &crate::net::SocketAddr, + stderr: &SharedWriter, +) -> crate::errors::Result { + let rustc = f.mk_bin("rustc").unwrap(); + let conn = connect_to_server(addr).unwrap(); + let cmdline = vec![ + "--crate-name".into(), + "dep".into(), + "src/lib.rs".into(), + "--emit=dep-info,metadata,link".into(), + ]; + let mut stdout = Cursor::new(Vec::new()); + let mut runtime = Runtime::new().unwrap(); + do_compile( + creator, + &mut runtime, + conn, + &rustc, + cmdline, + f.tempdir.path(), + Some(f.paths.clone()), + vec![], + &mut stdout, + &mut stderr.clone(), + ) +} + +#[test] +fn test_cli_rmeta_notification_delivery_from_daemon() { + use crate::compiler::ColorMode; + use crate::protocol::{CompileFinished, Response}; + use crate::util::write_length_prefixed_bincode; + + let _ = env_logger::try_init(); + let f = TestFixture::new(); + let stderr = SharedWriter::default(); + + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = crate::net::SocketAddr::Net(listener.local_addr().unwrap()); + let server_stderr = stderr.clone(); + let server = thread::spawn(move || -> Vec { + let mut sock = fake_daemon_until_rmeta_notification(listener); + + // This waits until the client writes what it has at this point, + // to ensure stderr gets the notification before we proceed to write more. + sock.set_read_timeout(Some(Duration::from_millis(5))) + .unwrap(); + let deadline = std::time::Instant::now() + Duration::from_secs(60); + while server_stderr.0.lock().unwrap().is_empty() { + match sock.peek(&mut [0]) { + Ok(0) => break, + Ok(_) => panic!("unexpected data from client"), + Err(e) if e.kind() == std::io::ErrorKind::ConnectionReset => break, + Err(_) => assert!( + std::time::Instant::now() < deadline, + "client neither wrote nor hung up" + ), + } + } + + let written_before_finished = server_stderr.0.lock().unwrap().clone(); + + // The client may have given up on us by now. + // The stderr is rustc's as-is, so it carries the notification too; + // the client is expected to drop that copy. + let _ = write_length_prefixed_bincode( + &mut sock, + Response::CompileFinished(CompileFinished { + retcode: Some(0), + signal: None, + stdout: vec![], + stderr: [RMETA_NOTIFICATION, OTHER_STDERR].concat(), + color_mode: ColorMode::Off, + }), + ); + written_before_finished + }); + + let retcode = compile_dep(new_creator(), &f, &addr, &stderr).unwrap(); + let written_before_finished = server.join().unwrap(); + + assert_eq!(0, retcode); + assert_eq!( + RMETA_NOTIFICATION, + written_before_finished.as_slice(), + "stderr written before CompileFinished" + ); + // The daemon leaves rustc's stderr intact. + // Client will dedup stderr if already forwarded. + assert_eq!( + [RMETA_NOTIFICATION, OTHER_STDERR].concat(), + *stderr.0.lock().unwrap(), + "stderr written in total" + ); +} + +/// This makes sure that if the daemon dies between the rmeta notification and codegen, +/// sccache falls back to a normal compilation. +#[test] +fn test_cli_rmeta_notification_delivery_after_daemon_disconnect() { + let _ = env_logger::try_init(); + let f = TestFixture::new(); + let stderr = SharedWriter::default(); + + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = crate::net::SocketAddr::Net(listener.local_addr().unwrap()); + let server = thread::spawn(move || { + let sock = fake_daemon_until_rmeta_notification(listener); + // The daemon dies mid-compile. + drop(sock); + }); + + // The fallback compile. + let creator = new_creator(); + next_command( + &creator, + Ok(MockChild::new( + exit_status(0), + "", + [RMETA_NOTIFICATION, OTHER_STDERR].concat(), + )), + ); + + let retcode = compile_dep(creator.clone(), &f, &addr, &stderr).unwrap(); + server.join().unwrap(); + + assert_eq!(0, retcode); + assert_eq!( + 0, + creator.lock().unwrap().children.len(), + "fallback rustc ran" + ); + // The fallback rustc emitted the notification the daemon had already streamed, + // and the client then dedups it. + assert_eq!( + [RMETA_NOTIFICATION, OTHER_STDERR].concat(), + *stderr.0.lock().unwrap(), + "stderr written in total" + ); +} diff --git a/src/util.rs b/src/util.rs index a2e752d23..081f6daad 100644 --- a/src/util.rs +++ b/src/util.rs @@ -460,15 +460,23 @@ pub fn fmt_duration_as_secs(duration: &Duration) -> String { format!("{}.{:03} s", duration.as_secs(), duration.subsec_millis()) } +/// Callback invoked with each complete stderr line (including its `\n`) as +/// soon as it is read, before the process has exited. +pub type StderrLineObserver = Box; + /// If `input`, write it to `child`'s stdin while also reading `child`'s stdout and stderr, then wait on `child` and return its status and output. /// /// This was lifted from `std::process::Child::wait_with_output` and modified /// to also write to stdin. -async fn wait_with_input_output(mut child: T, input: Option>) -> Result +async fn wait_with_input_output( + mut child: T, + input: Option>, + mut on_stderr_line: Option, +) -> Result where T: CommandChild + 'static, { - use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt}; let stdin = input.and_then(|i| { child.take_stdin().map(|mut stdin| async move { stdin.write_all(&i).await.context("failed to write stdin") @@ -494,10 +502,31 @@ where match stderr { Some(mut stderr) => { let mut buf = Vec::new(); - stderr - .read_to_end(&mut buf) - .await - .context("failed to read stderr")?; + match on_stderr_line.as_mut() { + None => { + stderr + .read_to_end(&mut buf) + .await + .context("failed to read stderr")?; + } + Some(observe) => { + let mut reader = tokio::io::BufReader::new(stderr); + loop { + let start = buf.len(); + // rustc's JSON emitter writes each notification followed by `b"\n"` + // so it should be safe to split. + // https://github.com/rust-lang/rust/blob/5ceaf6608eb354c2f5bbb3b8d974caa367dac81c/compiler/rustc_errors/src/json.rs#L88 + let read = reader + .read_until(b'\n', &mut buf) + .await + .context("failed to read stderr")?; + if read == 0 { + break; + } + observe(&buf[start..]); + } + } + } Result::Ok(Some(buf)) } None => Ok(None), @@ -526,7 +555,20 @@ where /// /// If the command returns a non-successful exit status, an error of `SccacheError::ProcessError` /// will be returned containing the process output. -pub async fn run_input_output(mut command: C, input: Option>) -> Result +pub async fn run_input_output(command: C, input: Option>) -> Result +where + C: RunCommand, +{ + run_input_output_observing(command, input, None).await +} + +/// Like [`run_input_output`] but additionally hands every complete stderr line +/// to `on_stderr_line` as it arrives. +pub async fn run_input_output_observing( + mut command: C, + input: Option>, + on_stderr_line: Option, +) -> Result where C: RunCommand, { @@ -541,7 +583,7 @@ where .spawn() .await?; - wait_with_input_output(child, input) + wait_with_input_output(child, input, on_stderr_line) .await .and_then(|output| { if output.status.success() { diff --git a/tests/sccache_rustc.rs b/tests/sccache_rustc.rs index 8fd4e2264..5dbe81fce 100644 --- a/tests/sccache_rustc.rs +++ b/tests/sccache_rustc.rs @@ -15,10 +15,12 @@ use std::{ path::{Path, PathBuf}, }; -struct StopServer; +struct StopServer(u16); impl Drop for StopServer { fn drop(&mut self) { let _ = Command::from_std(std::process::Command::new(env!("CARGO_BIN_EXE_sccache"))) + .env("SCCACHE_SERVER_PORT", self.0.to_string()) + .env_remove("SCCACHE_SERVER_UDS") .arg("--stop-server") .ok(); } @@ -58,19 +60,76 @@ fn test_symlinks() { let out_file = root.join("RUST_FILE"); symlink(root.join("rust1"), &rust).unwrap(); - drop(StopServer); - let _stop_server = StopServer; - run_sccache(root, &bin); + let port = 4321; + drop(StopServer(port)); + let _stop_server = StopServer(port); + run_sccache(root, &bin, port, false); let output1 = fs::read(&out_file).unwrap(); remove_file(&rust).unwrap(); symlink(root.join("rust2"), &rust).unwrap(); - run_sccache(root, &bin); + run_sccache(root, &bin, port, false); let output2 = fs::read(out_file).unwrap(); assert_ne!(output1, output2); } +#[test] +fn test_rmeta_notification_delivery_on_miss_and_hit() { + rmeta_notification_delivery_on_miss_and_hit(4322, false); +} + +#[test] +fn test_rmeta_notification_delivery_on_miss_and_hit_client_side() { + rmeta_notification_delivery_on_miss_and_hit(4323, true); +} + +/// Runs the real sccache and checks the caller sees the `.rmeta` +/// notification exactly once on both the miss and the hit. +fn rmeta_notification_delivery_on_miss_and_hit(port: u16, client_side: bool) { + let root = tempdir().unwrap(); + let root = root.path(); + + fs::write(root.join("counter"), b"0").unwrap(); + fs::write(root.join("RUST_FILE.rs"), []).unwrap(); + create_mock_rustc(root.join("rust")); + let bin = root.join("rust/bin"); + let out_file = root.join("RUST_FILE"); + + drop(StopServer(port)); + let _stop_server = StopServer(port); + + let miss = run_sccache(root, &bin, port, client_side); + let compiled_once = fs::read(&out_file).unwrap(); + let hit = run_sccache(root, &bin, port, client_side); + assert_eq!( + compiled_once, + fs::read(&out_file).unwrap(), + "second run was not a hit" + ); + + for (name, stderr) in [("miss", &miss.stderr), ("hit", &hit.stderr)] { + let stderr = String::from_utf8_lossy(stderr); + let notifications: Vec<&str> = stderr + .lines() + .filter(|line| line.contains("\"artifact\"")) + .collect(); + assert_eq!( + 1, + notifications + .iter() + .filter(|l| l.contains(".rmeta")) + .count(), + "{name}: .rmeta notification count in {stderr:?}" + ); + assert_eq!( + 1, + notifications.iter().filter(|l| l.contains(".rlib")).count(), + "{name}: .rlib notification count in {stderr:?}" + ); + } +} + fn create_mock_rustc(dir: PathBuf) { let bin = dir.join("bin"); create_dir_all(&bin).unwrap(); @@ -124,8 +183,10 @@ while [ "$#" -gt 0 ]; do done if [ "$build" -eq 1 ]; then + echo '{{"artifact":"'"$PWD"'/libsccache_rustc_tests.rmeta","emit":"metadata"}}' >&2 echo $(($(cat counter) + 1)) > counter cp counter RUST_FILE + echo '{{"artifact":"'"$PWD"'/libsccache_rustc_tests.rlib","emit":"link"}}' >&2 fi "#, dir.display(), @@ -137,7 +198,7 @@ fi set_permissions(&rustc, perm).unwrap(); } -fn run_sccache(root: &Path, path: &Path) { +fn run_sccache(root: &Path, path: &Path, port: u16, client_side: bool) -> std::process::Output { let mut paths: OsString = path.into(); paths.push(":"); paths.push(var_os("PATH").unwrap()); @@ -147,6 +208,9 @@ fn run_sccache(root: &Path, path: &Path) { .current_dir(root) .env("PATH", paths) .env("SCCACHE_DIR", root.join("sccache")) + .env("SCCACHE_SERVER_PORT", port.to_string()) + .env_remove("SCCACHE_SERVER_UDS") + .envs(client_side.then_some(("SCCACHE_CLIENT_SIDE", "1"))) .arg("rustc") .arg("RUST_FILE.rs") .arg("--crate-name=sccache_rustc_tests") @@ -154,5 +218,8 @@ fn run_sccache(root: &Path, path: &Path) { .arg("--emit=link") .arg("--out-dir") .arg(root) - .unwrap(); + .assert() + .success() + .get_output() + .clone() }