1use super::{PathInfo, PathInfoService};
2use crate::{
3 nar::{NarIngestionError, ingest_nar_and_hash},
4 pathinfoservice::{self, nix_http::castore_infused::try_infused_nar_path},
5};
6use futures::{TryStreamExt, stream::BoxStream};
7use nix_compat::{
8 narinfo::{self, NarInfo, Signature},
9 nixbase32,
10 nixhash::NixHash,
11};
12use reqwest::StatusCode;
13use snix_castore::{
14 blobservice::{self, BlobService},
15 directoryservice::{self, DirectoryService},
16 proto::{
17 blob_service_client::BlobServiceClient, directory_service_client::DirectoryServiceClient,
18 },
19};
20use snix_castore::{
21 composition::{CompositionContext, ServiceBuilder},
22 directoryservice::GRPCDirectoryService,
23};
24use std::sync::Arc;
25use tokio::io::{self, AsyncRead};
26use tonic::{async_trait, transport::Channel};
27use tracing::{Span, instrument, warn};
28use url::Url;
29
30mod castore_infused;
31
32pub struct NixHTTPPathInfoService<BS: Clone, DS> {
48 instance_name: String,
49 base_url: url::Url,
50 http_client: reqwest_middleware::ClientWithMiddleware,
51
52 blob_service: BS,
53 directory_service: DS,
54
55 layered_blob_service: blobservice::Cache<BS, blobservice::GRPCBlobService<Channel>>,
59 layered_directory_service:
63 directoryservice::combinators::Cache<DS, GRPCDirectoryService<Channel>>,
64
65 trusted_public_keys: Vec<narinfo::VerifyingKey>,
69
70 force_download_nar: bool,
72}
73
74impl<BS, DS> NixHTTPPathInfoService<BS, DS>
75where
76 BS: Clone,
77 DS: DirectoryService + Clone,
78{
79 pub fn try_build(
80 instance_name: String,
81 config: NixHTTPPathInfoServiceConfig,
82 blob_service: BS,
83 directory_service: DS,
84 ) -> Result<Self, Error> {
85 let mut trusted_public_keys = Vec::new();
86 for s in config.params.trusted_public_keys {
87 trusted_public_keys.push(
88 narinfo::VerifyingKey::parse(&s).map_err(|e| Error::ParseTrustedPublicKey(s, e))?,
89 )
90 }
91
92 let (layered_blob_service, layered_directory_service) = {
93 let grpc_url = {
94 let url_str = format!("grpc+{}", config.base_url);
95 let mut url: Url = url_str.parse().expect("url to parse");
96 url.set_path("");
97 url
98 };
99
100 let channel =
101 snix_castore::tonic::TonicConnector::from_url(&grpc_url)?.connect_expect_lazy();
102
103 let instance_name_layered = format!("{}-layered", instance_name);
104 let instance_name_grpc = format!("{}-grpc", instance_name);
105
106 (
107 blobservice::Cache::new(
108 instance_name_layered.clone(),
109 blob_service.clone(),
110 blobservice::GRPCBlobService::from_client(
111 instance_name_grpc.clone(),
112 BlobServiceClient::new(channel.clone()),
113 ),
114 ),
115 directoryservice::combinators::Cache::new(
116 instance_name_layered,
117 directory_service.clone(),
118 {
119 GRPCDirectoryService::from_client(
120 instance_name_grpc,
121 DirectoryServiceClient::new(channel).max_decoding_message_size(
122 directoryservice::GRPC_MAX_DECODING_MESSAGE_SIZE,
123 ),
124 )
125 },
126 ),
127 )
128 };
129
130 Ok(Self {
131 instance_name,
132 base_url: {
133 let mut base_url = config.base_url;
135 if !base_url.path().ends_with('/') {
136 let with_slash = format!("{}/", base_url.path());
137 base_url.set_path(&with_slash);
138 }
139 base_url
140 },
141 http_client: reqwest_middleware::ClientBuilder::new(
142 reqwest::Client::builder()
143 .user_agent(crate::USER_AGENT)
144 .build()
145 .map_err(reqwest_middleware::Error::Reqwest)?,
146 )
147 .with(snix_tracing::propagate::reqwest::tracing_middleware())
148 .build(),
149 blob_service,
150 directory_service,
151
152 layered_blob_service,
153 layered_directory_service,
154
155 trusted_public_keys,
156 force_download_nar: config.params.force_download_nar,
157 })
158 }
159
160 #[instrument(level=tracing::Level::TRACE, skip_all,fields(path.digest=nixbase32::encode(&digest)),err)]
161 fn derive_narinfo_url(&self, digest: [u8; 20]) -> Result<Url, Error> {
162 let s = format!("{}.narinfo", nixbase32::encode(&digest));
163 self.base_url
164 .join(&s)
165 .map_err(|e| Error::JoinUrl(self.base_url.to_owned(), s.to_owned(), e))
166 }
167}
168
169#[derive(Debug, thiserror::Error)]
170pub enum Error {
171 #[error("wrong arguments: {0}")]
172 WrongConfig(&'static str),
173 #[error("serde-qs error: {0}")]
174 SerdeQS(#[from] serde_qs::Error),
175 #[error("unable to parse pubkey {0}")]
176 ParseTrustedPublicKey(String, nix_compat::narinfo::VerifyingKeyError),
177 #[error("unable to construct tonic channel: {0}")]
178 TonicChannel(#[from] snix_castore::tonic::Error),
179
180 #[error("unable to join URL {0} with {1}")]
181 JoinUrl(Url, String, url::ParseError),
182 #[error("reqwest error")]
183 Reqwest(#[from] reqwest_middleware::Error),
184 #[error("unable to decode NARInfo response as string")]
185 DecodeBody(reqwest::Error),
186 #[error("unable to parse NARInfo")]
187 ParseNARInfo(nix_compat::narinfo::Error),
188 #[error("no valid signature found")]
189 NoValidSignature,
190 #[error("failed to request NAR, status {0}")]
191 FailedToRequestNAR(reqwest::StatusCode),
192 #[error("unsupported NAR compression: {0}")]
193 UnsupportedNARCompression(String),
194 #[error("failed to ingest NAR")]
195 IngestNAR(NarIngestionError),
196 #[error("NARSize mismatch, narinfo size {narinfo_size}, actual size {actual_size}")]
197 NARSizeMismatch { narinfo_size: u64, actual_size: u64 },
198 #[error("NARHash mismatch, narinfo NARHash {exp}, actual NARHash {act}",
199 exp = NixHash::Sha256(*.narinfo_nar_sha256),
200 act = NixHash::Sha256(*.actual_nar_sha256))]
201 NARHashMismatch {
202 narinfo_nar_sha256: [u8; 32],
203 actual_nar_sha256: [u8; 32],
204 },
205
206 #[error("put not supported")]
207 PutNotSupported,
208 #[error("list not supported")]
209 ListNotSupported,
210}
211
212#[async_trait]
213impl<BS, DS> PathInfoService for NixHTTPPathInfoService<BS, DS>
214where
215 BS: BlobService + Send + Sync + Clone + 'static,
216 DS: DirectoryService + Send + Sync + Clone + 'static,
217{
218 #[instrument(skip_all, err, fields(
219 path.digest=nixbase32::encode(&digest),
220 instance_name=%self.instance_name,
221 narinfo.url=tracing::field::Empty,
222 nar.url=tracing::field::Empty,
223 ))]
224 async fn get(&self, digest: [u8; 20]) -> Result<Option<PathInfo>, pathinfoservice::Error> {
225 let narinfo_url = self.derive_narinfo_url(digest)?;
226
227 let span = Span::current();
228 span.record("narinfo.url", narinfo_url.to_string());
229
230 let resp = self
231 .http_client
232 .get(narinfo_url)
233 .send()
234 .await
235 .map_err(Error::Reqwest)?;
236
237 if resp.status() == StatusCode::NOT_FOUND || resp.status() == StatusCode::FORBIDDEN {
241 return Ok(None);
242 }
243
244 let narinfo_str = resp.text().await.map_err(Error::DecodeBody)?;
245
246 let narinfo = NarInfo::parse(&narinfo_str).map_err(Error::ParseNARInfo)?;
248
249 if narinfo.store_path.digest() != &digest {
251 return Err("Store path digest in NARInfo doesn't match".into());
252 }
253
254 if !self.trusted_public_keys.is_empty() {
256 let fingerprint = narinfo.fingerprint();
257
258 if !self.trusted_public_keys.iter().any(|pubkey| {
259 narinfo
260 .signatures
261 .iter()
262 .any(|sig| pubkey.verify(&fingerprint, sig))
263 }) {
264 Err(Error::NoValidSignature)?
265 }
266 }
267
268 let root_node = if !self.force_download_nar
278 && let Some(root_node) = try_infused_nar_path(
279 &narinfo,
280 self.layered_blob_service.clone(),
281 &self.layered_directory_service,
282 )
283 .await
284 .unwrap_or_else(|err| {
285 warn!(%err, "unable to use infused store path");
286 None
287 }) {
288 root_node
289 } else {
290 let nar_url = self
292 .base_url
293 .join(narinfo.url)
294 .map_err(|e| Error::JoinUrl(self.base_url.clone(), narinfo.url.to_owned(), e))?;
295 span.record("nar.url", nar_url.to_string());
296
297 let resp = self
298 .http_client
299 .get(nar_url.clone())
300 .send()
301 .await
302 .map_err(Error::Reqwest)?;
303
304 if !resp.status().is_success() {
306 Err(Error::FailedToRequestNAR(resp.status()))?;
307 }
308
309 let r = tokio_util::io::StreamReader::new(resp.bytes_stream().map_err(|e| {
311 let e = e.without_url();
312 warn!(e=%e, "failed to get response body");
313 io::Error::new(io::ErrorKind::BrokenPipe, e.to_string())
314 }));
315
316 let mut r: Box<dyn AsyncRead + Send + Unpin> = match narinfo.compression {
318 None => Box::new(r) as Box<dyn AsyncRead + Send + Unpin>,
319 Some("bzip2") => Box::new(async_compression::tokio::bufread::BzDecoder::new(r))
320 as Box<dyn AsyncRead + Send + Unpin>,
321 Some("gzip") => Box::new(async_compression::tokio::bufread::GzipDecoder::new(r))
322 as Box<dyn AsyncRead + Send + Unpin>,
323 Some("xz") => Box::new(async_compression::tokio::bufread::XzDecoder::new(r))
324 as Box<dyn AsyncRead + Send + Unpin>,
325 Some("zstd") => {
326 let mut decoder = async_compression::tokio::bufread::ZstdDecoder::new(r);
329 decoder.multiple_members(true);
330 Box::new(decoder) as Box<dyn AsyncRead + Send + Unpin>
331 }
332 Some(comp_str) => Err(Error::UnsupportedNARCompression(comp_str.to_owned()))?,
333 };
334
335 let (root_node, nar_hash, nar_size) = ingest_nar_and_hash(
336 self.blob_service.clone(),
337 &self.directory_service,
338 &mut r,
339 &narinfo.ca,
340 )
341 .await
342 .map_err(Error::IngestNAR)?;
343
344 if narinfo.nar_size != nar_size {
346 Err(Error::NARSizeMismatch {
347 narinfo_size: narinfo.nar_size,
348 actual_size: nar_size,
349 })?
350 }
351 if narinfo.nar_hash != nar_hash {
352 Err(Error::NARHashMismatch {
353 narinfo_nar_sha256: narinfo.nar_hash,
354 actual_nar_sha256: nar_hash,
355 })?
356 }
357 root_node
358 };
359
360 Ok(Some(PathInfo {
361 store_path: narinfo.store_path.to_owned(),
362 node: root_node,
363 references: narinfo.references.iter().map(|sp| sp.to_owned()).collect(),
364 nar_size: narinfo.nar_size,
365 nar_sha256: narinfo.nar_hash,
366 deriver: narinfo.deriver.as_ref().map(|sp| sp.to_owned()),
367 signatures: narinfo
368 .signatures
369 .into_iter()
370 .map(|s| Signature::<String>::new(s.name().to_string(), s.bytes().to_owned()))
371 .collect(),
372 ca: narinfo.ca,
373 }))
374 }
375
376 #[instrument(skip_all, err, fields(
377 path.digest=nixbase32::encode(&digest),
378 instance_name=%self.instance_name,
379 narinfo.url=tracing::field::Empty,
380 ))]
381 async fn has(&self, digest: [u8; 20]) -> Result<bool, pathinfoservice::Error> {
382 let narinfo_url = self.derive_narinfo_url(digest)?;
383
384 let span = Span::current();
385 span.record("narinfo.url", narinfo_url.to_string());
386
387 let resp = self
388 .http_client
389 .head(narinfo_url)
390 .send()
391 .await
392 .map_err(Error::Reqwest)?;
393
394 if resp.status() == StatusCode::NOT_FOUND || resp.status() == StatusCode::FORBIDDEN {
398 Ok(false)
399 } else {
400 Ok(true)
401 }
402 }
403
404 #[instrument(skip_all, fields(path_info=?_path_info, instance_name=%self.instance_name))]
405 async fn put(&self, _path_info: PathInfo) -> Result<PathInfo, pathinfoservice::Error> {
406 Err(Box::new(Error::PutNotSupported))
407 }
408
409 fn list(&self) -> BoxStream<'static, Result<PathInfo, pathinfoservice::Error>> {
410 Box::pin(futures::stream::once(async {
411 Err(Error::ListNotSupported)?
412 }))
413 }
414}
415
416#[derive(serde::Deserialize, Clone, Debug, PartialEq, Eq)]
417#[serde(deny_unknown_fields)]
418pub struct NixHTTPPathInfoServiceConfig {
419 base_url: Url,
420
421 #[serde(flatten)]
422 params: NixHTTPPathInfoServiceParams,
423}
424
425#[derive(serde::Deserialize, Clone, Debug, PartialEq, Eq)]
426#[serde(deny_unknown_fields)]
427struct NixHTTPPathInfoServiceParams {
428 #[serde(default = "default_blob_service")]
429 blob_service: String,
430 #[serde(default = "default_directory_service")]
431 directory_service: String,
432 #[serde(default)]
433 trusted_public_keys: Vec<String>,
436
437 #[serde(default)]
438 force_download_nar: bool,
440}
441
442fn default_blob_service() -> String {
443 "&root".to_string()
444}
445fn default_directory_service() -> String {
446 "&root".to_string()
447}
448
449impl TryFrom<Url> for NixHTTPPathInfoServiceConfig {
450 type Error = Box<dyn std::error::Error + Send + Sync>;
451 fn try_from(url: Url) -> Result<Self, Self::Error> {
452 let scheme = url
453 .scheme()
454 .strip_prefix("nix+")
455 .ok_or_else(|| Error::WrongConfig("scheme must start with nix+"))?;
456
457 if !url.has_authority() {
458 Err(Error::WrongConfig("url must have authority component"))?
459 }
460 if !url.has_host() {
461 Err(Error::WrongConfig("url must have host component"))?
462 }
463 if !["http", "https"].contains(&scheme) {
464 Err(Error::WrongConfig("unknown scheme"))?
465 }
466
467 Ok(NixHTTPPathInfoServiceConfig {
468 base_url: {
474 let mut url: Url = url
475 .to_string()
476 .strip_prefix("nix+")
477 .unwrap()
478 .parse()
479 .expect("stripped URL to parse again");
480 url.set_query(None);
481 url
482 },
483 params: serde_qs::from_str(url.query().unwrap_or_default())?,
484 })
485 }
486}
487
488#[async_trait]
489impl ServiceBuilder for NixHTTPPathInfoServiceConfig {
490 type Output = dyn PathInfoService;
491 async fn build<'a>(
492 &'a self,
493 instance_name: &str,
494 context: &CompositionContext,
495 ) -> Result<Arc<Self::Output>, Box<dyn std::error::Error + Send + Sync + 'static>> {
496 let (blob_service, directory_service) = futures::join!(
497 context.resolve::<dyn BlobService>(&self.params.blob_service),
498 context.resolve::<dyn DirectoryService>(&self.params.directory_service)
499 );
500 let svc = NixHTTPPathInfoService::try_build(
501 instance_name.to_string(),
502 self.to_owned(),
503 blob_service?,
504 directory_service?,
505 )?;
506 Ok(Arc::new(svc))
507 }
508}
509
510#[cfg(test)]
511mod tests {
512 use super::{NixHTTPPathInfoServiceConfig, NixHTTPPathInfoServiceParams};
513 use rstest::rstest;
514 use url::Url;
515
516 #[rstest]
517 #[case::correct_nix_https("nix+https://cache.nixos.org", Some(
519 NixHTTPPathInfoServiceConfig {
520 base_url: "https://cache.nixos.org".try_into().unwrap(),
521 params: NixHTTPPathInfoServiceParams {
522 blob_service: "&root".to_string(),
523 directory_service: "&root".to_string(),
524 trusted_public_keys: vec![],
525 force_download_nar: false,
526 }
527 }
528 ))]
529 #[case::correct_nix_http("nix+http://cache.nixos.org", Some(
531 NixHTTPPathInfoServiceConfig {
532 base_url: "http://cache.nixos.org".try_into().unwrap(),
533 params: NixHTTPPathInfoServiceParams {
534 blob_service: "&root".to_string(),
535 directory_service: "&root".to_string(),
536 trusted_public_keys: vec![],
537 force_download_nar: false,
538 }
539 }
540 ))]
541 #[case::correct_nix_http_with_subpath("nix+http://192.0.2.1/foo", Some(
543 NixHTTPPathInfoServiceConfig {
544 base_url: "http://192.0.2.1/foo".try_into().unwrap(),
545 params: NixHTTPPathInfoServiceParams {
546 blob_service: "&root".to_string(),
547 directory_service: "&root".to_string(),
548 trusted_public_keys: vec![],
549 force_download_nar: false,
550 }
551 }
552 ))]
553 #[case::correct_nix_http_with_subpath_and_port("nix+http://[::1]:8080/foo", Some(
555 NixHTTPPathInfoServiceConfig {
556 base_url: "http://[::1]:8080/foo".try_into().unwrap(),
557 params: NixHTTPPathInfoServiceParams {
558 blob_service: "&root".to_string(),
559 directory_service: "&root".to_string(),
560 trusted_public_keys: vec![],
561 force_download_nar: false,
562 }
563 }
564
565 ))]
566 #[case::correct_nix_https_with_trusted_public_key(
568 "nix+https://cache.nixos.org?trusted_public_keys[0]=cache.nixos.org-1:6NCHdD59X431o0gWypbMrAURkbJ16ZPMQFGspcDShjY=", Some(
569 NixHTTPPathInfoServiceConfig {
570 base_url: "https://cache.nixos.org".try_into().unwrap(),
571 params: NixHTTPPathInfoServiceParams {
572 blob_service: "&root".to_string(),
573 directory_service: "&root".to_string(),
574 trusted_public_keys: vec![
575 "cache.nixos.org-1:6NCHdD59X431o0gWypbMrAURkbJ16ZPMQFGspcDShjY=".to_string()
576 ],
577 force_download_nar: false,
578 }
579 }
580 ))]
581 #[case::correct_nix_https_with_two_trusted_public_keys(
583 "nix+https://cache.nixos.org?trusted_public_keys[0]=cache.nixos.org-1:6NCHdD59X431o0gWypbMrAURkbJ16ZPMQFGspcDShjY=&trusted_public_keys[1]=foo:jp4fCEx9tBEId/L0ZsVJ26k0wC0fu7vJqLjjIGFkup8=", Some(
584 NixHTTPPathInfoServiceConfig {
585 base_url: "https://cache.nixos.org".try_into().unwrap(),
586 params: NixHTTPPathInfoServiceParams {
587 blob_service: "&root".to_string(),
588 directory_service: "&root".to_string(),
589 trusted_public_keys: vec![
590 "cache.nixos.org-1:6NCHdD59X431o0gWypbMrAURkbJ16ZPMQFGspcDShjY=".to_string(),
591 "foo:jp4fCEx9tBEId/L0ZsVJ26k0wC0fu7vJqLjjIGFkup8=".to_string()
592 ],
593 force_download_nar: false,
594 }
595 }
596 ))]
597 #[case::wrong_scheme("nix+grpc://example.com", None)]
598 #[case::missing_host("nix+http:///", None)]
599 #[case::missing_authority("nix+http:", None)]
600 #[case::trusted_public_keys_no_sequence(
602 "nix+https://cache.nixos.org?trusted_public_keys=cache.nixos.org-1:6NCHdD59X431o0gWypbMrAURkbJ16ZPMQFGspcDShjY=",
603 None
604 )]
605 #[case::trusted_public_keys_wrong_pubkey(
607 "nix+https://cache.nixos.org?trustedpublickeys=cache.nixos.org-1:6NCHdD59X431o0gWypbMrAURkbJ16ZPMQFGspcDShjY=",
608 None
609 )]
610 fn parse_url(#[case] url_str: &str, #[case] exp_config: Option<NixHTTPPathInfoServiceConfig>) {
611 let url: Url = url_str.parse().expect("url to parse");
612
613 match (NixHTTPPathInfoServiceConfig::try_from(url), exp_config) {
614 (Ok(_), None) => panic!("parsing url unexpectedly succeeded"),
615 (Ok(config), Some(exp_config)) => assert_eq!(exp_config, config),
616 (Err(_), None) => {}
617 (Err(e), Some(_)) => panic!("parsing url unexpectedly failed: {e}"),
618 }
619 }
620}