diff --git a/Cargo.lock b/Cargo.lock index 989b686349..deb46e5a48 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1490,8 +1490,8 @@ dependencies = [ [[package]] name = "harmonia-file-core" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "serde", "serde_json", @@ -1500,8 +1500,8 @@ dependencies = [ [[package]] name = "harmonia-file-nar" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "bstr", "bytes", @@ -1524,8 +1524,8 @@ dependencies = [ [[package]] name = "harmonia-protocol" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "async-stream", "bstr", @@ -1558,8 +1558,8 @@ dependencies = [ [[package]] name = "harmonia-protocol-derive" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "proc-macro2", "quote", @@ -1569,7 +1569,7 @@ dependencies = [ [[package]] name = "harmonia-store-aterm" version = "0.0.0-alpha.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "bytes", "harmonia-store-content-address", @@ -1584,8 +1584,8 @@ dependencies = [ [[package]] name = "harmonia-store-build-result" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "harmonia-store-derivation", "num_enum", @@ -1595,8 +1595,8 @@ dependencies = [ [[package]] name = "harmonia-store-content-address" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "derive_more", "harmonia-store-path", @@ -1608,7 +1608,7 @@ dependencies = [ [[package]] name = "harmonia-store-derivation" version = "0.0.0-alpha.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "bytes", "data-encoding", @@ -1626,8 +1626,8 @@ dependencies = [ [[package]] name = "harmonia-store-nar-info" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "harmonia-store-content-address", "harmonia-store-path", @@ -1640,8 +1640,8 @@ dependencies = [ [[package]] name = "harmonia-store-path" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "derive_more", "harmonia-utils-base-encoding", @@ -1653,8 +1653,8 @@ dependencies = [ [[package]] name = "harmonia-store-path-info" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "harmonia-store-content-address", "harmonia-store-path", @@ -1665,8 +1665,8 @@ dependencies = [ [[package]] name = "harmonia-store-remote" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "async-stream", "futures-core", @@ -1686,7 +1686,7 @@ dependencies = [ [[package]] name = "harmonia-utils-base-encoding" version = "0.0.0-alpha.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "data-encoding", "derive_more", @@ -1696,7 +1696,7 @@ dependencies = [ [[package]] name = "harmonia-utils-hash" version = "0.0.0-alpha.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "blake3", "data-encoding", @@ -1708,16 +1708,16 @@ dependencies = [ "sha1", "sha2 0.11.0", "thiserror", - "tokio", ] [[package]] name = "harmonia-utils-io" version = "0.0.0-alpha.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "bytes", "futures-util", + "harmonia-utils-hash", "libc", "nix", "pin-project-lite", @@ -1726,8 +1726,8 @@ dependencies = [ [[package]] name = "harmonia-utils-signature" -version = "3.2.0" -source = "git+https://github.com/nix-community/harmonia.git#b438172545da3af37998d0995b6440f1412ae14d" +version = "3.3.0" +source = "git+https://github.com/nix-community/harmonia.git#1424ae3f4e4218782f6ecd2714f989a7a1798bbc" dependencies = [ "data-encoding", "ed25519-dalek", @@ -1927,6 +1927,7 @@ dependencies = [ "harmonia-store-path", "harmonia-store-remote", "harmonia-utils-io", + "hydra-proto", "hydra-tracing", "listenfd", "sd-notify", @@ -1936,7 +1937,10 @@ dependencies = [ "test-utils", "thiserror", "tokio", + "tokio-stream", "toml 1.1.2+spec-1.1.0", + "tonic", + "tonic-prost", "tracing", ] diff --git a/subprojects/hydra-ad-hoc/Cargo.toml b/subprojects/hydra-ad-hoc/Cargo.toml index e865775434..364b58be11 100644 --- a/subprojects/hydra-ad-hoc/Cargo.toml +++ b/subprojects/hydra-ad-hoc/Cargo.toml @@ -22,7 +22,10 @@ serde = { workspace = true, features = [ "derive" ] } sqlx = { workspace = true, features = [ "runtime-tokio", "postgres" ] } thiserror.workspace = true tokio = { workspace = true, features = [ "full" ] } +tokio-stream.workspace = true toml.workspace = true +tonic.workspace = true +tonic-prost.workspace = true tracing.workspace = true db.workspace = true @@ -32,6 +35,8 @@ harmonia-store-derivation.workspace = true harmonia-store-path.workspace = true harmonia-store-remote.workspace = true harmonia-utils-io.workspace = true +hydra-proto = { workspace = true, features = [ "client" ] } + [features] otel = [ "hydra-tracing/otel" ] diff --git a/subprojects/hydra-ad-hoc/src/config.rs b/subprojects/hydra-ad-hoc/src/config.rs index 03be959314..e900d31bd8 100644 --- a/subprojects/hydra-ad-hoc/src/config.rs +++ b/subprojects/hydra-ad-hoc/src/config.rs @@ -97,6 +97,10 @@ fn default_data_dir() -> PathBuf { "/var/lib/hydra".into() } +fn default_queue_runner_grpc_addr() -> String { + "http://[::1]:50051".into() +} + /// Main configuration of the application #[derive(Debug, serde::Deserialize)] #[serde(deny_unknown_fields)] @@ -119,6 +123,11 @@ pub(crate) struct AppConfig { /// Hydra data directory; step logs are read from its `build-logs`. #[serde(default = "default_data_dir")] hydra_data_dir: PathBuf, + + /// `hydra-queue-runner`'s gRPC endpoint, used for the trusted-client + /// fast path (`SubmitDerivation`) that `build_derivation` takes. + #[serde(default = "default_queue_runner_grpc_addr")] + queue_runner_grpc_addr: String, } impl From for App { @@ -131,6 +140,7 @@ impl From for App { upstream_socket: val.upstream_socket, store_dir: val.store_dir, log_prefix: val.hydra_data_dir.join("build-logs"), + queue_runner_grpc_addr: val.queue_runner_grpc_addr, } } } @@ -142,6 +152,7 @@ pub(crate) struct App { pub upstream_socket: PathBuf, pub store_dir: StoreDir, pub log_prefix: PathBuf, + pub queue_runner_grpc_addr: String, } #[derive(Debug, thiserror::Error)] diff --git a/subprojects/hydra-ad-hoc/src/handler.rs b/subprojects/hydra-ad-hoc/src/handler.rs index 8481b8c7fa..797c0ba649 100644 --- a/subprojects/hydra-ad-hoc/src/handler.rs +++ b/subprojects/hydra-ad-hoc/src/handler.rs @@ -24,10 +24,11 @@ use harmonia_store_remote::pool::{ConnectionPool, PoolConfig}; use db::StoreDir; use db::models::{BuildID, BuildStatus}; +use hydra_proto::ad_hoc_service_client::AdHocServiceClient; use sqlx::Connection as _; use tokio::sync::{mpsc, oneshot, watch}; -use crate::logs::{LogSource, build_log_stream}; +use crate::logs::{LogSource, build_log_stream, inline_log_stream}; use crate::queries::{FinishedBuild, get_finished_build}; use crate::submit::{AdhocSubmitter, BuildRequest}; use crate::waiter::BuildWaiter; @@ -45,6 +46,9 @@ pub(crate) struct HydraDaemonHandler { waiter: BuildWaiter, submitter: AdhocSubmitter, logs: LogSource, + /// `hydra-queue-runner`'s `AdHocService`, used only by + /// `build_derivation`'s trusted-client fast path (see its doc comment). + queue_runner: AdHocServiceClient, } impl std::fmt::Debug for HydraDaemonHandler { @@ -63,6 +67,7 @@ impl HydraDaemonHandler { waiter: BuildWaiter, submitter: AdhocSubmitter, logs: LogSource, + queue_runner: AdHocServiceClient, ) -> Self { let upstream = ConnectionPool::with_store_dir( upstream_socket, @@ -76,6 +81,7 @@ impl HydraDaemonHandler { waiter, submitter, logs, + queue_runner, } } @@ -203,6 +209,70 @@ impl HydraDaemonHandler { result } + /// Dispatch `drv` on `hydra-queue-runner`'s `AdHocService` and wait for + /// the result. Used only by `build_derivation`'s trusted-client fast + /// path: there is no `.drv` on disk and no `Builds` row, so this + /// bypasses `schedule_build`/`run_build` (and everything they depend + /// on) entirely. + async fn run_inline_build( + &self, + drv_path: &StorePath, + drv: &BasicDerivation, + ) -> Result { + let drv_path_str = self.store_dir.display(drv_path).to_string(); + tracing::debug!(drv_path = %drv_path_str, "inline build requested"); + + let request = hydra_proto::SubmitDerivationRequest { + drv: Some(hydra_proto::nix::store::derivation::v1::Basic::from(drv)), + drv_path: Some(hydra_proto::ProtoStorePath::from(drv_path)), + wanted_outputs: Vec::new(), + }; + let mut stream = self + .queue_runner + .clone() + .submit_derivation(request) + .await + .map_err(|e| ProtocolError::custom(format!("submit_derivation: {e}")))? + .into_inner(); + + // First event is always `Dispatched`; only logged, the client + // learns about progress from the log stream this build's tail + // (`inline_log_stream`) is following independently. + match stream + .message() + .await + .map_err(|e| ProtocolError::custom(format!("submit_derivation: {e}")))? + { + Some(hydra_proto::SubmitDerivationEvent { + event: Some(hydra_proto::submit_derivation_event::Event::Dispatched(d)), + }) => { + tracing::info!(drv_path = %drv_path_str, machine = %d.machine_hostname, "inline build dispatched"); + } + _ => { + return Err(ProtocolError::custom( + "submit_derivation: expected a Dispatched event first", + )); + } + } + + let result = match stream + .message() + .await + .map_err(|e| ProtocolError::custom(format!("submit_derivation: {e}")))? + { + Some(hydra_proto::SubmitDerivationEvent { + event: Some(hydra_proto::submit_derivation_event::Event::Result(r)), + }) => r, + _ => { + return Err(ProtocolError::custom( + "submit_derivation: stream ended before a Result event", + )); + } + }; + + inline_result_to_build_result(&drv_path_str, &result) + } + /// Run `work` while streaming the logs of every build it announces. /// /// The daemon protocol delivers an operation's log messages before @@ -344,8 +414,55 @@ fn finished_to_build_result( }) } +/// Build a [`BuildResult`] straight from the queue runner's +/// `SubmitDerivationResult` — there is no `Builds`/`buildoutputs` row for +/// an inline build, so this bypasses `get_finished_build`/`queries.rs` +/// entirely, unlike [`finished_to_build_result`]. +fn inline_result_to_build_result( + drv_path: &str, + result: &hydra_proto::SubmitDerivationResult, +) -> Result { + let inner = if result.success { + let mut built_outputs = BTreeMap::new(); + for (name, path) in &result.outputs { + let name: OutputName = name + .parse() + .map_err(|e| ProtocolError::custom(format!("invalid output name: {e}")))?; + built_outputs.insert( + name, + UnkeyedRealisation { + out_path: path.0.clone(), + signatures: BTreeSet::new(), + }, + ); + } + if built_outputs.is_empty() { + return Err(ProtocolError::custom(format!( + "build of {drv_path} succeeded but reported no outputs" + ))); + } + BuildResultInner::Success(BuildResultSuccess { + status: SuccessStatus::Built, + built_outputs, + }) + } else { + BuildResultInner::Failure(BuildResultFailure { + status: FailureStatus::PermanentFailure, + error_msg: result.error_msg.clone().into(), + is_non_deterministic: false, + }) + }; + Ok(BuildResult { + inner, + times_built: 1, + start_time: 0, + stop_time: 0, + cpu_user: None, + cpu_system: None, + }) +} + /// Synthesize realisations from recorded output paths; missing paths are queue-runner bugs. -/// /// The result is keyed by output name, so each entry only carries the /// unkeyed half of the realisation — the `DrvOutput` key is implied by the /// derivation being built. @@ -417,8 +534,23 @@ impl HandshakeDaemonStore for HydraDaemonHandler { } impl DaemonStore for HydraDaemonHandler { + /// `NotTrusted`, not `Trusted`: this is what a connecting client sees + /// during the handshake (`remoteTrustsUs` in Nix's own terms), and it + /// decides which code path Nix's `ssh-ng://`/`ssh://` build-hook + /// (`build-remote.cc`) takes for a delegated build. A `Trusted` daemon + /// makes the hook send input-addressed derivations inline via + /// `BuildDerivation` (no `.drv` ever uploaded, no `Builds` row to hang + /// the build off — see `build_derivation`'s `run_inline_build` path). + /// `NotTrusted` makes it upload the real `.drv` first and call + /// `build_paths` instead, so those builds get Hydra's normal + /// `Steps`/`Queues` graph: dedup, previous-failure caching, and + /// wake-on-completion across the whole build. Content-addressed + /// derivations are unaffected either way — Nix always takes the inline + /// fast path for those, trusted or not, since CA output paths are + /// self-verifying by content hash — so `build_derivation`'s inline path + /// stays load-bearing for that case. fn trust_level(&self) -> Option { - Some(TrustLevel::Trusted) + Some(TrustLevel::NotTrusted) } fn set_options<'a>( @@ -429,6 +561,15 @@ impl DaemonStore for HydraDaemonHandler { ready(Ok(())).empty_logs() } + /// Dispatch a derivation inline: Nix's trusted-client build-hook fast + /// path never uploads a `.drv` for it (it sends the resolved derivation + /// content directly in the daemon protocol's `BuildDerivation` call), + /// so there is no `.drv` on disk for the queue runner's normal + /// ingestion pipeline to read and no `Builds` row to hang the build off + /// (see `hydra-queue-runner`'s `State::dispatch_inline_derivation`, + /// which this calls over gRPC). This bypasses `AdhocSubmitter`, + /// `BuildWaiter` and `queries::get_finished_build` entirely: those all + /// key off a `Builds` row this build never gets. fn build_derivation<'a>( &'a mut self, drv_path: &'a StorePath, @@ -438,18 +579,28 @@ impl DaemonStore for HydraDaemonHandler { let this = self.clone(); let drv_path = drv_path.clone(); let drv = drv.clone(); - self.with_live_logs(move |announce| async move { - require_normal_mode(mode)?; - this.assert_drv_uploaded(&drv_path).await?; - let drv_path_str = this.store_dir.display(&drv_path).to_string(); - let nix_name: String = drv.name.to_string(); - let system = std::str::from_utf8(&drv.platform) - .map_err(|e| ProtocolError::custom(format!("non-utf8 platform: {e}")))?; - let finished = this - .run_build(&drv_path_str, &nix_name, system, &announce) - .await?; - finished_to_build_result(&drv_path_str, &finished) - }) + let (done_tx, done_rx) = watch::channel(false); + let logs_src = this.logs.clone(); + let log_drv_path = drv_path.clone(); + let (result_tx, result_rx) = oneshot::channel(); + tokio::spawn(async move { + let result = async { + require_normal_mode(mode)?; + this.run_inline_build(&drv_path, &drv).await + } + .await; + let _ = done_tx.send(true); + let _ = result_tx.send(result); + }); + let logs = inline_log_stream(logs_src, log_drv_path, done_rx); + async move { + result_rx.await.unwrap_or_else(|_| { + Err(ProtocolError::custom( + "build task ended without reporting a result", + )) + }) + } + .with_logs(logs) } fn build_paths<'a>( diff --git a/subprojects/hydra-ad-hoc/src/logs.rs b/subprojects/hydra-ad-hoc/src/logs.rs index 881631f59d..9a270cd893 100644 --- a/subprojects/hydra-ad-hoc/src/logs.rs +++ b/subprojects/hydra-ad-hoc/src/logs.rs @@ -154,13 +154,45 @@ fn step_log_stream( src: LogSource, ev: StepEvent, ids: Arc, - mut finished: watch::Receiver, + finished: watch::Receiver, ) -> impl Stream + Send { stream! { let Some(drv) = lookup_drv(&src, ev).await else { return; }; let id = ids.fetch_add(1, Ordering::Relaxed); + let inner = follow_drv_log(src, drv, id, finished); + let mut inner = Box::pin(inner); + while let Some(msg) = inner.next().await { + yield msg; + } + } +} + +/// A single build's worth of messages for a derivation that is not backed +/// by a `BuildSteps` row: a `Build` activity, its log lines as +/// `BuildLogLine` results, and the activity's end. Used by +/// [`crate::handler::HydraDaemonHandler::build_derivation`]'s inline +/// dispatch path, where the drv path (and hence the log file) is already +/// known and there is no step announcement to look it up from. +pub(crate) fn inline_log_stream( + src: LogSource, + drv: harmonia_store_path::StorePath, + finished: watch::Receiver, +) -> impl Stream + Send { + follow_drv_log(src, drv, 1, finished) +} + +/// Tail `drv`'s log file into `Build`/`BuildLogLine`/end-of-activity +/// messages, from when this is called until `finished` turns true and the +/// grace period after it elapses. +fn follow_drv_log( + src: LogSource, + drv: harmonia_store_path::StorePath, + id: u64, + mut finished: watch::Receiver, +) -> impl Stream + Send { + stream! { let drv_str = src.store_dir.display(&drv).to_string(); yield LogMessage::StartActivity(Activity { // Nix's own `Build` activity: derivation, machine, round, rounds. @@ -178,10 +210,10 @@ fn step_log_stream( }); let path = build_logs::log_path(&src.log_prefix, &drv); - tracing::debug!(?ev, path = %path.display(), "step started; waiting for its log"); + tracing::debug!(path = %path.display(), "step started; waiting for its log"); if wait_for_log_file(&path, &mut finished).await { let mut sub = src.tails.subscribe(&path).await; - tracing::debug!(?ev, backlog = sub.backlog.len(), finished = *finished.borrow(), "following step log"); + tracing::debug!(backlog = sub.backlog.len(), finished = *finished.borrow(), "following step log"); // `step_finished` does not mean the log is complete: the builder // sends its log and its result on separate streams, so the queue // runner can still be writing the file. Close the tail a grace @@ -231,7 +263,7 @@ fn step_log_stream( } } - tracing::debug!(?ev, "step log done"); + tracing::debug!("step log done"); yield LogMessage::StopActivity(StopActivity { id }); } } diff --git a/subprojects/hydra-ad-hoc/src/main.rs b/subprojects/hydra-ad-hoc/src/main.rs index 86e7af8684..5f8551778f 100644 --- a/subprojects/hydra-ad-hoc/src/main.rs +++ b/subprojects/hydra-ad-hoc/src/main.rs @@ -27,6 +27,7 @@ mod config; mod handler; mod logs; mod queries; +mod queue_runner; mod server; mod submit; mod waiter; @@ -55,6 +56,7 @@ async fn main() -> eyre::Result<()> { db::Database::new(config.db_url.expose_secret(), config.max_db_connections).await?; let waiter = BuildWaiter::start(&database).await?; let submitter = AdhocSubmitter::new(database.clone()).await?; + let queue_runner = queue_runner::connect(&config.queue_runner_grpc_addr)?; let logs = LogSource { db: database.clone(), store_dir: store_dir.clone(), @@ -70,6 +72,7 @@ async fn main() -> eyre::Result<()> { waiter, submitter, logs, + queue_runner, ); let server = match &cli.socket { diff --git a/subprojects/hydra-ad-hoc/src/queue_runner.rs b/subprojects/hydra-ad-hoc/src/queue_runner.rs new file mode 100644 index 0000000000..a4c62360f7 --- /dev/null +++ b/subprojects/hydra-ad-hoc/src/queue_runner.rs @@ -0,0 +1,27 @@ +//! gRPC client to `hydra-queue-runner`'s `AdHocService`, used by +//! [`crate::handler::HydraDaemonHandler::build_derivation`] for Nix's +//! trusted-client fast path (see the module docs on `handler.rs`). +//! +//! This client has no TLS configuration of its own. If the queue runner's +//! gRPC listener has mTLS enabled (`state.mtls.enabled()`, see +//! `hydra-queue-runner/src/server/grpc.rs`), that listener requires every +//! connection — this one included — to present a client certificate; a +//! plain `http://` connection will fail outright rather than merely being +//! treated as unauthenticated. `hydra-builder` already carries the +//! `ClientTlsConfig`/`Identity`/`Certificate` wiring this would need +//! (`hydra-builder/src/config.rs` and `grpc.rs:198-260`) — this client +//! needs the same treatment before `hydra-ad-hoc` can be deployed against +//! an mTLS-enabled queue runner. Tracked as a known gap, not yet fixed. + +use hydra_proto::ad_hoc_service_client::AdHocServiceClient; + +/// A lazy channel: connecting does not block on the queue runner being up +/// (or on this call), matching the "start even if the queue runner isn't +/// there yet" tolerance the rest of `hydra-ad-hoc` has for its Postgres +/// dependency. +pub(crate) fn connect( + addr: &str, +) -> Result, tonic::transport::Error> { + let endpoint = tonic::transport::Endpoint::from_shared(addr.to_owned())?; + Ok(AdHocServiceClient::new(endpoint.connect_lazy())) +} diff --git a/subprojects/hydra-manual/src/architecture.md b/subprojects/hydra-manual/src/architecture.md index 8f6afd3371..673b423849 100644 --- a/subprojects/hydra-manual/src/architecture.md +++ b/subprojects/hydra-manual/src/architecture.md @@ -78,6 +78,7 @@ graph BT hydra-proto --> nix-support hydra-ad-hoc --> build-logs hydra-ad-hoc --> db + hydra-ad-hoc --> hydra-proto hydra-ad-hoc --> hydra-tracing hydra-evaluator --> db hydra-evaluator --> hydra-tracing diff --git a/subprojects/hydra-queue-runner/src/server/adhoc_grpc.rs b/subprojects/hydra-queue-runner/src/server/adhoc_grpc.rs new file mode 100644 index 0000000000..895827f566 --- /dev/null +++ b/subprojects/hydra-queue-runner/src/server/adhoc_grpc.rs @@ -0,0 +1,133 @@ +//! `AdHocService`: dispatches derivations `hydra-ad-hoc` received inline +//! over the daemon protocol, bypassing `Builds`/`Steps`/`Queues` entirely. +//! See [`crate::state::State::dispatch_inline_derivation`] for why: Nix's +//! trusted-client build-hook fast path never uploads a `.drv` file for +//! these, so there is nothing for the normal ingestion pipeline to read. +//! +//! Unlike [`super::grpc::RunnerService`], this service is called by +//! `hydra-ad-hoc`, not by builder machines, so it is not wrapped in +//! [`super::grpc::CheckAuthInterceptor`] and has no builder-facing auth story +//! of its own yet (see where it is added to the server in `grpc.rs`). + +use std::sync::Arc; + +use harmonia_store_derivation::derivation::BasicDerivation; +use harmonia_store_derivation::derived_path::OutputName; + +use hydra_proto::ad_hoc_service_server::AdHocService; +use hydra_proto::{ + Dispatched, SubmitDerivationEvent, SubmitDerivationRequest, submit_derivation_event, +}; + +use crate::state::State; + +type AdHocResult = Result, tonic::Status>; +type SubmitDerivationResponseStream = std::pin::Pin< + Box> + Send>, +>; + +#[allow(missing_debug_implementations)] +#[derive(Clone)] +pub struct Server { + state: Arc, +} + +impl Server { + #[must_use] + pub const fn new(state: Arc) -> Self { + Self { state } + } +} + +#[tonic::async_trait] +impl AdHocService for Server { + type SubmitDerivationStream = SubmitDerivationResponseStream; + + #[tracing::instrument(skip(self, req), err)] + async fn submit_derivation( + &self, + req: tonic::Request, + ) -> AdHocResult { + let req = req.into_inner(); + let drv: BasicDerivation = req + .drv + .ok_or_else(|| tonic::Status::invalid_argument("missing drv"))? + .try_into() + .map_err(|e| tonic::Status::invalid_argument(format!("invalid drv: {e}")))?; + let drv_path = req + .drv_path + .ok_or_else(|| tonic::Status::invalid_argument("missing drv_path"))? + .0; + let wanted_outputs = req + .wanted_outputs + .into_iter() + .map(|s| { + s.parse::().map_err(|e| { + tonic::Status::invalid_argument(format!("invalid output name: {e}")) + }) + }) + .collect::, _>>()?; + + let (internal_build_id, machine_hostname, rx) = self + .state + .dispatch_inline_derivation(drv, drv_path, wanted_outputs) + .await + .map_err(|e| tonic::Status::unavailable(format!("dispatch failed: {e}")))?; + + let (tx, out_rx) = tokio::sync::mpsc::channel(2); + tokio::spawn(async move { + if tx + .send(Ok(SubmitDerivationEvent { + event: Some(submit_derivation_event::Event::Dispatched(Dispatched { + machine_hostname, + })), + })) + .await + .is_err() + { + return; + } + + let result = match rx.await { + Ok(result) => build_result_to_event(&result), + Err(_) => SubmitDerivationEvent { + event: Some(submit_derivation_event::Event::Result( + hydra_proto::SubmitDerivationResult { + success: false, + error_msg: format!( + "builder for {internal_build_id} disappeared before reporting a result" + ), + outputs: std::collections::HashMap::new(), + }, + )), + }, + }; + let _ = tx.send(Ok(result)).await; + }); + + Ok(tonic::Response::new( + Box::pin(tokio_stream::wrappers::ReceiverStream::new(out_rx)) + as Self::SubmitDerivationStream, + )) + } +} + +/// Translate the builder's `BuildResultInfo` into the event +/// `hydra-ad-hoc` gets over `SubmitDerivation`. +fn build_result_to_event(result: &hydra_proto::BuildResultInfo) -> SubmitDerivationEvent { + let success = result.result_state() == hydra_proto::BuildResultState::Success; + let mut outputs = std::collections::HashMap::with_capacity(result.output_infos.len()); + for (name, info) in &result.output_infos { + let Some(path) = &info.path else { continue }; + outputs.insert(name.clone(), path.clone()); + } + SubmitDerivationEvent { + event: Some(submit_derivation_event::Event::Result( + hydra_proto::SubmitDerivationResult { + success, + error_msg: result.error_msg.clone().unwrap_or_default(), + outputs, + }, + )), + } +} diff --git a/subprojects/hydra-queue-runner/src/server/grpc.rs b/subprojects/hydra-queue-runner/src/server/grpc.rs index a8167f6572..88dc71e428 100644 --- a/subprojects/hydra-queue-runner/src/server/grpc.rs +++ b/subprojects/hydra-queue-runner/src/server/grpc.rs @@ -216,10 +216,24 @@ impl Server { .build_v1()?; let (_health_reporter, health_service) = tonic_health::server::health_reporter(); + // `AdHocServiceServer` is called by hydra-ad-hoc, not builder + // machines, so it does not go through `CheckAuthInterceptor` (that + // interceptor's token auth is a builder-only auth story). When mTLS + // is configured, the `tls_config` above already requires every + // service on this listener -- this one included -- to present a + // client certificate signed by `client_ca_cert`; without mTLS this + // service is unauthenticated, the same accepted-risk posture as the + // health/reflection services added alongside it here. + let adhoc_service = hydra_proto::ad_hoc_service_server::AdHocServiceServer::new( + crate::server::adhoc_grpc::Server::new(state.clone()), + ) + .max_decoding_message_size(50 * 1024 * 1024) + .max_encoding_message_size(50 * 1024 * 1024); server .add_service(health_service) .add_service(reflection_service) .add_service(intercepted_service) + .add_service(adhoc_service) .serve_with_incoming(incoming) .await?; @@ -506,6 +520,14 @@ impl RunnerService for Server { tonic::Status::invalid_argument("machine_id is not a valid uuid.") })?; + // Builds dispatched by `dispatch_inline_derivation` (hydra-ad-hoc's + // trusted-client fast path) have no `Builds`/`Steps` row to + // finalize; hand the result straight to whoever is waiting on it + // instead of falling into the normal succeed/fail-by-uuid path. + if state.resolve_inline_completion(build_id, machine_id, req.clone()) { + return Ok(tonic::Response::new(hydra_proto::Empty {})); + } + // Finalize inline and propagate failure: remove_job must run to free // the machine slot, so a failed finalize has to reach the builder // (which retries complete_build) rather than being swallowed. diff --git a/subprojects/hydra-queue-runner/src/server/mod.rs b/subprojects/hydra-queue-runner/src/server/mod.rs index 8eb7894859..1e8063c782 100644 --- a/subprojects/hydra-queue-runner/src/server/mod.rs +++ b/subprojects/hydra-queue-runner/src/server/mod.rs @@ -1,2 +1,3 @@ +pub mod adhoc_grpc; pub mod grpc; pub mod http; diff --git a/subprojects/hydra-queue-runner/src/state/mod.rs b/subprojects/hydra-queue-runner/src/state/mod.rs index c28389cbe3..fea12bade8 100644 --- a/subprojects/hydra-queue-runner/src/state/mod.rs +++ b/subprojects/hydra-queue-runner/src/state/mod.rs @@ -82,6 +82,8 @@ pub enum StateLogicError { MachineLookup(#[from] MachineLookupError), #[error(transparent)] DrvLookup(#[from] DrvLookupError), + #[error(transparent)] + InlineDispatch(#[from] InlineDispatchError), } impl From for StateError { @@ -374,6 +376,16 @@ pub struct State { /// Cached from `nix_daemon_config.real_store_dir()`. /// `None` means the logical store dir is the filesystem path. pub real_store_dir: Option, + + /// Completion channels for builds dispatched by + /// [`State::dispatch_inline_derivation`]: these have no `Builds`/`Steps` + /// row, so `complete_build` cannot finalize them via + /// `succeed_step_by_uuid`/`fail_step_by_uuid`. Keyed by the job's + /// `internal_build_id`, resolved and removed once by whichever + /// `complete_build` call reports that id. + inline_completions: parking_lot::Mutex< + HashMap>, + >, } impl State { @@ -491,6 +503,7 @@ impl State { upload_completion_rx: parking_lot::Mutex::new(Some(upload_completion_rx)), real_store_dir: nix_config.real_store_dir(), nix_daemon_config: nix_config, + inline_completions: parking_lot::Mutex::new(HashMap::new()), config, })) } @@ -2352,7 +2365,125 @@ impl State { self.fail_step(machine_id, &drv_path, state, timings, error_msg) .await } +} + +/// Errors dispatching a derivation received inline over the daemon +/// protocol (see [`State::dispatch_inline_derivation`]). +#[derive(Debug, thiserror::Error)] +pub enum InlineDispatchError { + #[error("non-utf8 platform")] + InvalidPlatformUtf8(#[from] std::str::Utf8Error), + + #[error("no machine available for system '{0}'")] + NoMachineForSystem(String), +} + +impl From for StateError { + fn from(e: InlineDispatchError) -> Self { + Self::Logic(StateLogicError::InlineDispatch(e)) + } +} + +impl State { + /// Dispatch a derivation whose content arrived inline (no `.drv` file on + /// disk to read): pick a machine directly from [`Machines`], bypassing + /// `Builds`/`Steps`/`Queues` entirely, and hand it to the machine's + /// `build_drv` the same way `realise_drv_on_valid_machine` does. + /// + /// There is no `Builds` row backing this build, so it gets none of the + /// GC-root pinning or priority propagation a queued build gets: the + /// caller (`hydra-ad-hoc`) is on its own for keeping outputs alive. The + /// completion is delivered on the returned channel once a builder's + /// `complete_build` reports the minted id (see `Server::complete_build` + /// and `State::resolve_inline_completion`). + #[tracing::instrument(skip(self, drv), err)] + pub async fn dispatch_inline_derivation( + &self, + drv: harmonia_store_derivation::derivation::BasicDerivation, + drv_path: StorePath, + wanted_outputs: Vec, + ) -> Result< + ( + uuid::Uuid, + String, + tokio::sync::oneshot::Receiver, + ), + StateError, + > { + let _ = &wanted_outputs; // Nix always builds every output; kept for future filtering. + let system = std::str::from_utf8(&drv.platform) + .map_err(InlineDispatchError::from)? + .to_owned(); + let required_features = step::required_features(&drv); + + let machine = self + .machines + .get_machine_for_system( + &system, + &required_features, + Some(self.config.get_machine_free_fn()), + ) + .ok_or_else(|| InlineDispatchError::NoMachineForSystem(system.clone()))?; + + let internal_build_id = uuid::Uuid::new_v4(); + let (tx, rx) = tokio::sync::oneshot::channel(); + self.inline_completions.lock().insert(internal_build_id, tx); + + let mut job = machine::Job::new(BuildID::MAX, drv_path.clone()); + job.internal_build_id = internal_build_id; + job.result.set_start_time_now(); + + let resolved_drv = hydra_proto::nix::store::derivation::v1::Basic::from(&drv); + + if let Err(e) = machine + .build_drv( + job, + drv_path, + self.config.max_log_size(), + self.config.max_output_size(), + self.config.max_silent_time(), + self.config.build_timeout(), + None, + resolved_drv, + ) + .await + { + self.inline_completions.lock().remove(&internal_build_id); + return Err(e.into()); + } + + Ok((internal_build_id, machine.hostname.clone(), rx)) + } + /// If `build_id` is a job dispatched by [`State::dispatch_inline_derivation`], + /// remove its machine slot and hand the reported result to whoever is + /// waiting on the returned receiver, and report `true` so the caller + /// (`complete_build`) skips the normal `Steps`/`Queues` finalization path. + /// Returns `false` for a `build_id` that is not an inline job (the normal + /// path is not this build's). + #[tracing::instrument(skip(self, result), fields(%build_id, %machine_id))] + pub fn resolve_inline_completion( + &self, + build_id: uuid::Uuid, + machine_id: uuid::Uuid, + result: hydra_proto::BuildResultInfo, + ) -> bool { + let Some(tx) = self.inline_completions.lock().remove(&build_id) else { + return false; + }; + if let Some(machine) = self.machines.get_machine_by_id(machine_id) + && let Some(job) = machine.get_job_drv_for_build_id(build_id) + { + machine.remove_job(&job); + } + // The receiver may already be gone (e.g. the ad-hoc client + // disconnected mid-build); the builder's job is done either way. + let _ = tx.send(result); + true + } +} + +impl State { #[allow(clippy::too_many_lines)] #[tracing::instrument(skip(self, machine, job, step), fields(%drv_path), err)] async fn inner_fail_job( diff --git a/subprojects/hydra-queue-runner/src/state/step.rs b/subprojects/hydra-queue-runner/src/state/step.rs index d720660615..97729d1b17 100644 --- a/subprojects/hydra-queue-runner/src/state/step.rs +++ b/subprojects/hydra-queue-runner/src/state/step.rs @@ -129,7 +129,15 @@ pub(crate) struct StepDrvInfo { /// Such derivations carry the attribute only inside the `__json` blob (parsed /// into `structured_attrs`) with no flat env var, so reading the env alone /// mis-schedules them onto machines lacking the feature (e.g. `big-parallel`). -fn required_features(drv: &Derivation) -> Vec { +/// +/// Generic over `Inputs`/`Output` so it works on both the full [`Derivation`] +/// (parsed `.drv` with unresolved `Built` inputs) and a force-resolved +/// [`harmonia_store_derivation::derivation::BasicDerivation`] received +/// inline (see `crate::state::State::dispatch_inline_derivation`): both are +/// the same `DerivationT` shape and this only reads `env`/`structured_attrs`. +pub(crate) fn required_features( + drv: &harmonia_store_derivation::derivation::DerivationT, +) -> Vec { if let Some(structured) = &drv.structured_attrs { let Some(serde_json::Value::Array(features)) = structured.attrs.get("requiredSystemFeatures") @@ -889,6 +897,29 @@ mod tests { assert_eq!(required_features(&drv), vec!["big-parallel"]); } + /// `dispatch_inline_derivation` (the trusted-client fast path) calls + /// `required_features` on a `BasicDerivation`, not a parsed + /// `Derivation` -- the generic `DerivationT` signature + /// must actually work on that type, not just happen to compile against + /// it via unused generic parameters. + #[test] + fn required_features_on_basic_derivation() { + let drv = harmonia_store_derivation::derivation::BasicDerivation { + name: "test".parse().unwrap(), + outputs: BTreeMap::new(), + inputs: harmonia_store_path::StorePathSet::new(), + platform: bytes::Bytes::from_static(b"x86_64-linux"), + builder: bytes::Bytes::from_static(b"/bin/sh"), + args: Vec::new(), + env: BTreeMap::from([( + bytes::Bytes::from_static(b"requiredSystemFeatures"), + bytes::Bytes::from_static(b"kvm"), + )]), + structured_attrs: None, + }; + assert_eq!(required_features(&drv), vec!["kvm"]); + } + #[test] fn pending_runnable_drained_once() { let steps = Steps::new(); diff --git a/subprojects/proto/v1/streaming.proto b/subprojects/proto/v1/streaming.proto index 1fd26f5bd7..ca62465eae 100644 --- a/subprojects/proto/v1/streaming.proto +++ b/subprojects/proto/v1/streaming.proto @@ -25,6 +25,39 @@ service RunnerService { rpc NotifyPresignedUploadComplete(PresignedUploadComplete) returns (Empty) {} } +// Called by hydra-ad-hoc (not builders) to dispatch a derivation whose +// content arrived inline over the daemon protocol: Nix's trusted-client +// build-hook fast path never uploads a `.drv` file for it, so there is no +// `Builds`/`Steps`/`Queues` row to hang the build off. The queue runner +// picks a machine and dispatches directly; see +// `State::dispatch_inline_derivation`. +service AdHocService { + rpc SubmitDerivation(SubmitDerivationRequest) returns (stream SubmitDerivationEvent) {} +} + +message SubmitDerivationRequest { + nix.store.derivation.v1.Basic drv = 1; + nix.store.v1.StorePath drv_path = 2; + repeated string wanted_outputs = 3; +} + +message SubmitDerivationEvent { + oneof event { + Dispatched dispatched = 1; + SubmitDerivationResult result = 2; + } +} + +message Dispatched { + string machine_hostname = 1; +} + +message SubmitDerivationResult { + bool success = 1; + string error_msg = 2; + map outputs = 3; +} + message Empty {} message VersionCheckRequest {