mod get_package_error;
mod node_package_client;
mod node_package_information;
use crate::{
common::copy_dir_recursive, config::rt::RtcBuild,
pipelines::node_packages::node_package_client::NodePackageClient,
};
use anyhow::Context;
use async_compression::tokio::bufread::GzipDecoder;
use futures_util::{StreamExt, TryStreamExt, stream::FuturesUnordered};
use std::{io, sync::Arc};
use tokio::{fs::remove_dir_all, task::JoinHandle};
use tokio_util::io::StreamReader;
pub type NodePackageHandles = FuturesUnordered<JoinHandle<anyhow::Result<()>>>;
pub fn spawn_node_packages(cfg: Arc<RtcBuild>) -> NodePackageHandles {
tracing::info!("node packages {:?}", cfg.node_packages);
let working_directory = cfg.working_directory.clone();
let futures: FuturesUnordered<_> = cfg
.node_packages
.iter()
.map(|node_package_cfg| {
let package_information = format!(
"{}@{}{}",
node_package_cfg.name,
node_package_cfg.version,
node_package_cfg
.registry
.clone()
.map(|registry| format!("(registry: {registry})"))
.unwrap_or_default()
);
let node_package_cfg = node_package_cfg.clone();
let working_directory = working_directory.clone();
tokio::spawn(async move {
let http_node_module_client = if let Some(registry) = node_package_cfg.registry {
NodePackageClient::new(®istry)?
} else {
NodePackageClient::default()
};
let target_path = node_package_cfg.target_path.unwrap_or(format!(
"target/node_modules/{}/{}",
node_package_cfg.name, node_package_cfg.version
));
let target_path = working_directory.join(target_path);
if target_path.exists() {
tracing::debug!(
"target path ({}) already exists, skipping",
target_path.display()
);
return Ok(());
}
tracing::info!("download node package {package_information}");
let package = http_node_module_client
.get_package(&node_package_cfg.name, &node_package_cfg.version)
.await
.with_context(|| {
format!("failed to retrieve Node package: {package_information}")
})?;
let tarball = reqwest::get(package.distribution.tarball)
.await?
.bytes_stream();
let tarball = tarball.map_err(io::Error::other);
let archive_data = GzipDecoder::new(StreamReader::new(tarball));
let archive = async_tar::Archive::new(archive_data);
archive.unpack(&target_path).await?;
let package_directory = target_path.join("package");
tracing::debug!("move from {package_directory:?} to {target_path:?}");
copy_dir_recursive(package_directory.clone(), target_path).await?;
remove_dir_all(package_directory).await?;
tracing::info!("finished to download node package {package_information}");
Ok(())
})
})
.collect();
futures
}
pub async fn wait_node_packages(mut futures: NodePackageHandles) -> anyhow::Result<()> {
while let Some(result) = futures.next().await {
result??;
}
Ok(())
}