use std::collections::BTreeMap;
use std::path::PathBuf;
use std::time::Duration;
use rowan::TextRange;
use crate::cache::Cache;
use crate::diag::{DiagCode, Diagnostic, FileId, Severity, Span};
use crate::lock::{LockEntry, Lockfile};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FetchError {
pub message: String,
}
const FETCH_TIMEOUT: Duration = Duration::from_secs(30);
pub fn fetch(url: &str) -> Result<Vec<u8>, FetchError> {
fetch_with_timeout(url, FETCH_TIMEOUT)
}
fn fetch_with_timeout(url: &str, timeout: Duration) -> Result<Vec<u8>, FetchError> {
let agent: ureq::Agent = ureq::Agent::config_builder()
.timeout_global(Some(timeout))
.build()
.into();
let mut response = agent.get(url).call().map_err(|err| FetchError {
message: err.to_string(),
})?;
response.body_mut().read_to_vec().map_err(|err| FetchError {
message: err.to_string(),
})
}
fn is_fetchable_url(url: &str) -> bool {
if url.chars().any(|c| c.is_whitespace() || c.is_control()) {
return false;
}
let Some(rest) = url
.strip_prefix("https://")
.or_else(|| url.strip_prefix("http://"))
else {
return false;
};
let host_end = rest.find(['/', '?', '#']).unwrap_or(rest.len());
!rest[..host_end].is_empty()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Frozen {
Yes,
No,
}
impl From<bool> for Frozen {
fn from(frozen: bool) -> Self {
if frozen { Frozen::Yes } else { Frozen::No }
}
}
pub fn materialize_imports(
imports: &BTreeMap<String, String>,
lock: Option<&Lockfile>,
cache: &Cache,
frozen: Frozen,
) -> (Vec<(String, PathBuf)>, Lockfile, Vec<Diagnostic>) {
let mut resolved = Vec::new();
let mut regenerated = Lockfile::default();
let mut diagnostics = Vec::new();
let mut urls: Vec<&String> = imports.values().collect();
urls.sort();
urls.dedup();
for url in urls {
if !is_fetchable_url(url) {
diagnostics.push(detached(
DiagCode::MANI_101,
format!("cannot fetch `{url}`: not a fetchable `http(s)` URL"),
));
continue;
}
let pinned = lock
.and_then(|lock| lock.entries.get(url))
.map(|entry| entry.sha256.clone());
if let Some(sha) = &pinned
&& let Some(path) = cache.lookup(url, sha)
{
regenerated.entries.insert(
url.clone(),
LockEntry {
sha256: sha.clone(),
},
);
resolved.push((url.clone(), path));
continue;
}
match frozen {
Frozen::Yes => match pinned {
None => diagnostics.push(detached(
DiagCode::MANI_103,
format!("`--frozen`: no lockfile entry for `{url}`"),
)),
Some(_) => diagnostics.push(detached(
DiagCode::MANI_104,
format!("`--frozen`: `{url}` is pinned but not in the cache"),
)),
},
Frozen::No => {
let bytes = match fetch(url) {
Ok(bytes) => bytes,
Err(err) => {
diagnostics.push(detached(
DiagCode::MANI_101,
format!("failed to fetch `{url}`: {}", err.message),
));
continue;
}
};
let (sha, path) = match cache.store(url, &bytes) {
Ok(stored) => stored,
Err(err) => {
diagnostics.push(detached(
DiagCode::MANI_101,
format!("failed to cache `{url}`: {err}"),
));
continue;
}
};
if let Some(expected) = &pinned
&& expected != &sha
{
diagnostics.push(detached(
DiagCode::MANI_102,
format!(
"hash mismatch for `{url}`: the lockfile pins `{expected}` but the fetched content hashes to `{sha}`"
),
));
continue;
}
regenerated
.entries
.insert(url.clone(), LockEntry { sha256: sha });
resolved.push((url.clone(), path));
}
}
}
(resolved, regenerated, diagnostics)
}
fn detached(code: DiagCode, message: String) -> Diagnostic {
Diagnostic {
code,
severity: Severity::Error,
message,
primary: Span {
file: FileId::DETACHED,
range: TextRange::default(),
},
labels: Vec::new(),
fixits: Vec::new(),
}
}
#[cfg(test)]
mod tests {
use std::io::{BufRead as _, BufReader, Write as _};
use std::net::{TcpListener, TcpStream};
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::thread::JoinHandle;
use super::*;
struct TempDir(PathBuf);
impl TempDir {
fn new(label: &str) -> Self {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let mut path = std::env::temp_dir();
path.push(format!(
"ridl-core-fetch-{label}-{}-{}",
std::process::id(),
COUNTER.fetch_add(1, Ordering::SeqCst),
));
std::fs::create_dir_all(&path).expect("create the temp dir");
Self(path)
}
fn path(&self) -> &Path {
&self.0
}
}
impl Drop for TempDir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
fn make_tar(name: &str, contents: &[u8]) -> Vec<u8> {
let mut builder = tar::Builder::new(Vec::new());
let mut header = tar::Header::new_gnu();
header.set_size(contents.len() as u64);
header.set_mode(0o644);
header.set_cksum();
builder
.append_data(&mut header, name, contents)
.expect("append the file to the tar");
builder.into_inner().expect("finish the tar")
}
struct Stub {
url: String,
addr: std::net::SocketAddr,
hits: Arc<AtomicUsize>,
shutdown: Arc<AtomicBool>,
handle: Option<JoinHandle<()>>,
}
impl Stub {
fn serving(body: Vec<u8>) -> Stub {
let listener = TcpListener::bind("127.0.0.1:0").expect("bind the stub");
let addr = listener.local_addr().expect("stub address");
let hits = Arc::new(AtomicUsize::new(0));
let shutdown = Arc::new(AtomicBool::new(false));
let handle = {
let hits = hits.clone();
let shutdown = shutdown.clone();
std::thread::spawn(move || {
for incoming in listener.incoming() {
if shutdown.load(Ordering::SeqCst) {
break;
}
let stream = match incoming {
Ok(stream) => stream,
Err(_) => break,
};
if shutdown.load(Ordering::SeqCst) {
break;
}
serve_one(stream, &body, &hits);
}
})
};
Stub {
url: format!("http://{addr}/package.tar"),
addr,
hits,
shutdown,
handle: Some(handle),
}
}
fn url(&self) -> &str {
&self.url
}
fn hits(&self) -> usize {
self.hits.load(Ordering::SeqCst)
}
fn reset_hits(&self) {
self.hits.store(0, Ordering::SeqCst);
}
}
impl Drop for Stub {
fn drop(&mut self) {
self.shutdown.store(true, Ordering::SeqCst);
let _ = TcpStream::connect(self.addr);
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}
fn serve_one(mut stream: TcpStream, body: &[u8], hits: &AtomicUsize) {
let mut reader = BufReader::new(stream.try_clone().expect("clone the stub stream"));
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line) {
Ok(0) => return, Ok(_) if line == "\r\n" || line == "\n" => break,
Ok(_) => {}
Err(_) => return,
}
}
hits.fetch_add(1, Ordering::SeqCst);
let header = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/x-tar\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
body.len(),
);
let _ = stream.write_all(header.as_bytes());
let _ = stream.write_all(body);
let _ = stream.flush();
}
struct StallingStub {
addr: std::net::SocketAddr,
shutdown: Arc<AtomicBool>,
handle: Option<JoinHandle<()>>,
}
impl StallingStub {
fn start() -> StallingStub {
let listener = TcpListener::bind("127.0.0.1:0").expect("bind the stalling stub");
let addr = listener.local_addr().expect("stub address");
let shutdown = Arc::new(AtomicBool::new(false));
let handle = {
let shutdown = shutdown.clone();
std::thread::spawn(move || {
let mut held = Vec::new();
for incoming in listener.incoming() {
if shutdown.load(Ordering::SeqCst) {
break;
}
match incoming {
Ok(stream) => held.push(stream),
Err(_) => break,
}
}
drop(held);
})
};
StallingStub {
addr,
shutdown,
handle: Some(handle),
}
}
fn url(&self) -> String {
format!("http://{}/package.tar", self.addr)
}
}
impl Drop for StallingStub {
fn drop(&mut self) {
self.shutdown.store(true, Ordering::SeqCst);
let _ = TcpStream::connect(self.addr);
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}
fn cache(dir: &TempDir) -> Cache {
Cache {
root: dir.path().join("cache"),
}
}
fn imports_of(url: &str) -> BTreeMap<String, String> {
BTreeMap::from([("remote.pkg".to_string(), url.to_string())])
}
#[test]
fn fetch_returns_the_served_bytes() {
let body = make_tar("a.typl", b"package remote.pkg\n");
let stub = Stub::serving(body.clone());
let got = fetch(stub.url()).expect("the fetch succeeds");
assert_eq!(got, body, "the fetched bytes are the served body");
assert_eq!(stub.hits(), 1, "exactly one request reached the stub");
}
#[test]
fn fetch_times_out_on_a_stalled_server() {
let stub = StallingStub::start();
let start = std::time::Instant::now();
let result = fetch_with_timeout(&stub.url(), Duration::from_millis(300));
assert!(result.is_err(), "a stalled server must time out, not hang");
assert!(
start.elapsed() < Duration::from_secs(5),
"the timeout must fire promptly (no hang), took {:?}",
start.elapsed(),
);
}
#[test]
fn fetch_of_a_dead_address_errors() {
let listener = TcpListener::bind("127.0.0.1:0").expect("bind");
let addr = listener.local_addr().expect("addr");
drop(listener);
let url = format!("http://{addr}/package.tar");
assert!(fetch(&url).is_err(), "a refused connection is a FetchError");
}
#[test]
fn materialize_unpacks_a_remote_package_into_the_cache() {
let source = "package remote.pkg\ntype Speed: km/h\n";
let stub = Stub::serving(make_tar("remote.typl", source.as_bytes()));
let dir = TempDir::new("materialize");
let cache = cache(&dir);
let imports = imports_of(stub.url());
let (resolved, lock, diags) = materialize_imports(&imports, None, &cache, Frozen::No);
assert!(diags.is_empty(), "a clean fetch, no diagnostics: {diags:?}");
assert_eq!(stub.hits(), 1, "the artifact is fetched exactly once");
assert_eq!(resolved.len(), 1);
let (resolved_url, path) = &resolved[0];
assert_eq!(resolved_url, stub.url());
let unpacked = std::fs::read_to_string(path.join("remote.typl"))
.expect("the .typl file is unpacked into the cache");
assert_eq!(unpacked, source);
let entry = lock.entries.get(stub.url()).expect("the URL is pinned");
assert_eq!(entry.sha256.len(), 64, "the pin is a SHA-256 hex string");
}
#[test]
fn a_cache_hit_skips_the_fetch() {
let stub = Stub::serving(make_tar("remote.typl", b"package remote.pkg\n"));
let dir = TempDir::new("cache-hit");
let cache = cache(&dir);
let imports = imports_of(stub.url());
let (_, lock, diags) = materialize_imports(&imports, None, &cache, Frozen::No);
assert!(diags.is_empty());
assert_eq!(stub.hits(), 1, "the first run fetches once");
stub.reset_hits();
let (resolved, lock2, diags2) =
materialize_imports(&imports, Some(&lock), &cache, Frozen::No);
assert!(diags2.is_empty());
assert_eq!(stub.hits(), 0, "a cache hit makes no request");
assert_eq!(
resolved.len(),
1,
"the package still resolves from the cache"
);
assert_eq!(lock2, lock, "the regenerated lockfile keeps the same pin");
}
#[test]
fn a_hash_mismatch_is_mani_102() {
let stub = Stub::serving(make_tar("remote.typl", b"package remote.pkg\n"));
let dir = TempDir::new("mismatch");
let cache = cache(&dir); let imports = imports_of(stub.url());
let mut lock = Lockfile::default();
lock.entries.insert(
stub.url().to_string(),
LockEntry {
sha256: "0".repeat(64),
},
);
let (resolved, regenerated, diags) =
materialize_imports(&imports, Some(&lock), &cache, Frozen::No);
assert_eq!(stub.hits(), 1, "the mismatch is found only after fetching");
assert_eq!(codes(&diags), vec!["MANI-102"]);
assert!(resolved.is_empty(), "a mismatched import does not resolve");
assert!(
regenerated.entries.is_empty(),
"a mismatched hash is never pinned",
);
}
#[test]
fn frozen_success_uses_the_cache_without_fetching() {
let stub = Stub::serving(make_tar("remote.typl", b"package remote.pkg\n"));
let dir = TempDir::new("frozen-ok");
let cache = cache(&dir);
let imports = imports_of(stub.url());
let (_, lock, _) = materialize_imports(&imports, None, &cache, Frozen::No);
stub.reset_hits();
let (resolved, _, diags) = materialize_imports(&imports, Some(&lock), &cache, Frozen::Yes);
assert!(
diags.is_empty(),
"a cached pin resolves cleanly under --frozen"
);
assert_eq!(stub.hits(), 0, "--frozen never fetches");
assert_eq!(resolved.len(), 1);
}
#[test]
fn frozen_missing_lockfile_entry_is_mani_103() {
let stub = Stub::serving(make_tar("remote.typl", b"package remote.pkg\n"));
let dir = TempDir::new("frozen-103");
let cache = cache(&dir);
let imports = imports_of(stub.url());
let (resolved, _, diags) = materialize_imports(&imports, None, &cache, Frozen::Yes);
assert_eq!(codes(&diags), vec!["MANI-103"]);
assert_eq!(stub.hits(), 0, "--frozen never fetches");
assert!(resolved.is_empty());
}
#[test]
fn frozen_pinned_but_uncached_is_mani_104() {
let stub = Stub::serving(make_tar("remote.typl", b"package remote.pkg\n"));
let dir = TempDir::new("frozen-104");
let cache = cache(&dir); let imports = imports_of(stub.url());
let mut lock = Lockfile::default();
lock.entries.insert(
stub.url().to_string(),
LockEntry {
sha256: "a".repeat(64),
},
);
let (resolved, _, diags) = materialize_imports(&imports, Some(&lock), &cache, Frozen::Yes);
assert_eq!(codes(&diags), vec!["MANI-104"]);
assert_eq!(stub.hits(), 0, "--frozen never fetches");
assert!(resolved.is_empty());
}
#[test]
fn an_invalid_url_is_mani_101_without_fetching() {
let dir = TempDir::new("bad-url");
let cache = cache(&dir);
let imports = BTreeMap::from([("bad.dep".to_string(), "not-a-url".to_string())]);
let (resolved, _, diags) = materialize_imports(&imports, None, &cache, Frozen::No);
assert_eq!(codes(&diags), vec!["MANI-101"]);
assert!(resolved.is_empty());
}
fn codes(diags: &[Diagnostic]) -> Vec<&str> {
diags.iter().map(|d| d.code.as_str()).collect()
}
}