From 0f005df305dfa283f505242d3c3a027f4723fb0b Mon Sep 17 00:00:00 2001 From: tbereknyei Date: Fri, 18 Sep 2026 15:55:05 -0400 Subject: [PATCH 1/3] hydra-ad-hoc: dispatch trusted-client fast-path builds inline Nix's ssh-ng/ssh build-hook trusted-client fast path sends BuildDerivation with the derivation content inline and never uploads a .drv file, so hydra-ad-hoc's previous Builds-row-and-wait flow could never load it and the client hung forever. Add an AdHocService.SubmitDerivation RPC that lets hydra-queue-runner dispatch such a derivation straight to a machine, bypassing Steps/Queues/Builds entirely, with a minimal UUID-keyed completion path back to hydra-ad-hoc. Also flip hydra-ad-hoc's trust_level() to NotTrusted: this makes Nix take the .drv-uploading build_paths path for input-addressed derivations instead of the inline fast path, so those get Hydra's normal dedup, previous-failure caching, and wake-on-completion for free. Content-addressed derivations still use the new inline path, since Nix takes that fast path for CA derivations regardless of trust. Known gaps: the new hydra-ad-hoc -> hydra-queue-runner gRPC client has no TLS config, so it cannot connect when the queue runner has mTLS enabled; and plain ssh:// (the legacy protocol) always behaves as trusted with no way to signal otherwise, so it keeps using the inline path even for input-addressed derivations. Co-Authored-By: Claude Sonnet 5 --- Cargo.lock | 4 + subprojects/hydra-ad-hoc/Cargo.toml | 5 + subprojects/hydra-ad-hoc/src/config.rs | 11 ++ subprojects/hydra-ad-hoc/src/handler.rs | 181 ++++++++++++++++-- subprojects/hydra-ad-hoc/src/logs.rs | 40 +++- subprojects/hydra-ad-hoc/src/main.rs | 3 + subprojects/hydra-ad-hoc/src/queue_runner.rs | 27 +++ .../src/server/adhoc_grpc.rs | 133 +++++++++++++ .../hydra-queue-runner/src/server/grpc.rs | 22 +++ .../hydra-queue-runner/src/server/mod.rs | 1 + .../hydra-queue-runner/src/state/mod.rs | 131 +++++++++++++ .../hydra-queue-runner/src/state/step.rs | 33 +++- subprojects/proto/v1/streaming.proto | 33 ++++ 13 files changed, 604 insertions(+), 20 deletions(-) create mode 100644 subprojects/hydra-ad-hoc/src/queue_runner.rs create mode 100644 subprojects/hydra-queue-runner/src/server/adhoc_grpc.rs diff --git a/Cargo.lock b/Cargo.lock index 989b686349..87d59401f7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -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..edf66fee57 100644 --- a/subprojects/hydra-ad-hoc/Cargo.toml +++ b/subprojects/hydra-ad-hoc/Cargo.toml @@ -22,6 +22,9 @@ serde = { workspace = true, features = [ "derive" ] } sqlx = { workspace = true, features = [ "runtime-tokio", "postgres" ] } thiserror.workspace = true tokio = { workspace = true, features = [ "full" ] } +tokio-stream.workspace = true +tonic.workspace = true +tonic-prost.workspace = true toml.workspace = true tracing.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-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 { From 53203c7b00fa6f3961ba0514ee2118e166c2ee24 Mon Sep 17 00:00:00 2001 From: tbereknyei Date: Fri, 18 Sep 2026 20:07:33 -0400 Subject: [PATCH 2/3] deps: bump harmonia to pick up fixed-output derivation hash fix hydra-ad-hoc's build_derivation handler (the ssh-ng/ssh trusted-client fast path) failed on any fixed-output derivation with "invalid nixbase32 hash at position N": harmonia's read_derivation_output/ write_derivation_output hardcoded Base::NixBase32 for a CAFixed output's hash, but real Nix's derivation::write/parseDerivation (aterm.cc) encode that hash as hex (Base16) -- the same ATerm writer is shared between .drv files on disk and the BuildDerivation wire RPC, and both use hex there. Already fixed upstream in nix-community/harmonia#1177 (merged 2026-08-27, e6a597b): the write side now uses Base::Hex, and the read side auto-detects the encoding by length instead of assuming one fixed base, matching how Nix's own hash parsing behaves. Bumping the pin picks this up with no source changes needed on our side. Co-Authored-By: Claude Sonnet 5 --- Cargo.lock | 56 +++++++++++++++++++++++++++--------------------------- 1 file changed, 28 insertions(+), 28 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 87d59401f7..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", From 6917f4fb18f30c4e91b792f4ddb3ca3215aeb326 Mon Sep 17 00:00:00 2001 From: tbereknyei Date: Sat, 19 Sep 2026 07:11:43 -0400 Subject: [PATCH 3/3] fix CI: treefmt Cargo.toml ordering, regenerate dependency diagram - subprojects/hydra-ad-hoc/Cargo.toml: toml.workspace was alphabetically misplaced after tonic-prost instead of before tonic. - subprojects/hydra-manual/src/architecture.md: regenerate to include the new hydra-ad-hoc -> hydra-proto edge added by the queue-runner client this branch introduces. Verified with `nix flake check` (all checks green) after both fixes. Co-Authored-By: Claude Sonnet 5 --- subprojects/hydra-ad-hoc/Cargo.toml | 2 +- subprojects/hydra-manual/src/architecture.md | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/subprojects/hydra-ad-hoc/Cargo.toml b/subprojects/hydra-ad-hoc/Cargo.toml index edf66fee57..364b58be11 100644 --- a/subprojects/hydra-ad-hoc/Cargo.toml +++ b/subprojects/hydra-ad-hoc/Cargo.toml @@ -23,9 +23,9 @@ sqlx = { workspace = true, features = [ "runtime-tokio", "post thiserror.workspace = true tokio = { workspace = true, features = [ "full" ] } tokio-stream.workspace = true +toml.workspace = true tonic.workspace = true tonic-prost.workspace = true -toml.workspace = true tracing.workspace = true db.workspace = true 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