use crate::AssetSyncProgressRenderer;
use crate::asset::config::AssetConfig;
use crate::asset::content::Content;
use crate::asset::content_encoder::ContentEncoder;
use crate::batch_upload::semaphores::Semaphores;
use crate::canister_api::methods::chunk::create_chunk;
use crate::canister_api::methods::chunk::create_chunks;
use crate::canister_api::types::asset::AssetDetails;
use crate::error::CreateChunkError;
use crate::error::CreateEncodingError;
use crate::error::CreateEncodingError::EncodeContentFailed;
use crate::error::CreateProjectAssetError;
use crate::error::SetEncodingError;
use candid::Nat;
use futures::TryFutureExt;
use futures::future::try_join_all;
use ic_utils::Canister;
use mime::Mime;
use slog::{Logger, debug};
use std::collections::BTreeMap;
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::sync::Mutex;
const CONTENT_ENCODING_IDENTITY: &str = "identity";
const MAX_COST_SINGLE_FILE_MB: usize = 45;
pub(crate) const MAX_CHUNK_SIZE: usize = 1_900_000;
#[derive(Clone, Debug)]
pub(crate) struct AssetDescriptor {
pub(crate) source: PathBuf,
pub(crate) key: String,
pub(crate) config: AssetConfig,
}
pub(crate) struct ProjectAssetEncoding {
pub(crate) uploader_chunk_ids: Vec<usize>,
pub(crate) sha256: Vec<u8>,
pub(crate) already_in_place: bool,
}
pub(crate) struct ProjectAsset {
pub(crate) asset_descriptor: AssetDescriptor,
pub(crate) media_type: Mime,
pub(crate) encodings: HashMap<String, ProjectAssetEncoding>,
}
#[derive(Debug, PartialEq)]
pub enum Mode {
ByProposal,
NormalDeploy,
}
type IdMapping = BTreeMap<usize, Nat>;
type UploadQueue = Vec<(usize, Vec<u8>)>;
type CanisterChunkSizeMap = HashMap<Nat, usize>;
pub(crate) struct ChunkUploader<'agent> {
canister: Canister<'agent>,
batch_id: Nat,
api_version: u16,
chunks: Arc<AtomicUsize>,
bytes: Arc<AtomicUsize>,
id_mapping: Arc<Mutex<IdMapping>>,
upload_queue: Arc<Mutex<UploadQueue>>,
canister_chunk_sizes: Arc<Mutex<CanisterChunkSizeMap>>,
}
impl<'agent> ChunkUploader<'agent> {
pub(crate) fn new(canister: Canister<'agent>, api_version: u16, batch_id: Nat) -> Self {
Self {
canister,
batch_id,
api_version,
chunks: Arc::new(AtomicUsize::new(0)),
bytes: Arc::new(AtomicUsize::new(0)),
id_mapping: Arc::new(Mutex::new(BTreeMap::new())),
upload_queue: Arc::new(Mutex::new(vec![])),
canister_chunk_sizes: Arc::new(Mutex::new(HashMap::new())),
}
}
pub(crate) async fn create_chunk(
&self,
contents: &[u8],
semaphores: &Semaphores,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<usize, CreateChunkError> {
let uploader_chunk_id = self.chunks.fetch_add(1, Ordering::SeqCst);
let chunk_size = contents.len();
self.bytes.fetch_add(chunk_size, Ordering::SeqCst);
if chunk_size == MAX_CHUNK_SIZE || self.api_version < 2 {
let canister_chunk_id = create_chunk(
&self.canister,
&self.batch_id,
contents,
semaphores,
progress,
)
.await?;
self.id_mapping
.lock()
.await
.insert(uploader_chunk_id, canister_chunk_id.clone());
self.canister_chunk_sizes
.lock()
.await
.insert(canister_chunk_id, chunk_size);
Ok(uploader_chunk_id)
} else {
self.add_to_upload_queue(uploader_chunk_id, contents).await;
self.upload_chunks(4 * MAX_CHUNK_SIZE, semaphores, progress)
.await?;
Ok(uploader_chunk_id)
}
}
pub(crate) async fn finalize_upload(
&self,
semaphores: &Semaphores,
mode: Mode,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<(), CreateChunkError> {
let max_retained_bytes = if mode == Mode::ByProposal {
0
} else {
MAX_CHUNK_SIZE / 2
};
self.upload_chunks(max_retained_bytes, semaphores, progress)
.await
}
pub(crate) fn bytes(&self) -> usize {
self.bytes.load(Ordering::SeqCst)
}
pub(crate) fn chunks(&self) -> usize {
self.chunks.load(Ordering::SeqCst)
}
pub(crate) async fn uploader_ids_to_canister_chunk_ids(
&self,
uploader_ids: &[usize],
) -> Result<(Vec<Nat>, Option<Vec<u8>>), SetEncodingError> {
let mut chunk_ids = vec![];
let mut last_chunk: Option<Vec<u8>> = None;
let mapping = self.id_mapping.lock().await;
let queue = self.upload_queue.lock().await;
for uploader_id in uploader_ids {
if let Some(item) = mapping.get(uploader_id) {
chunk_ids.push(item.clone());
} else if let Some(last_chunk_data) = queue
.iter()
.find_map(|(id, data)| if id == uploader_id { Some(data) } else { None })
{
match last_chunk.as_mut() {
Some(existing_data) => existing_data.extend(last_chunk_data.iter()),
None => last_chunk = Some(last_chunk_data.clone()),
}
} else {
return Err(SetEncodingError::UnknownUploaderChunkId(*uploader_id));
}
}
Ok((chunk_ids, last_chunk))
}
async fn add_to_upload_queue(&self, uploader_chunk_id: usize, contents: &[u8]) {
let mut queue = self.upload_queue.lock().await;
queue.push((uploader_chunk_id, contents.into()));
}
async fn upload_chunks(
&self,
max_retained_bytes: usize,
semaphores: &Semaphores,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<(), CreateChunkError> {
let mut queue = self.upload_queue.lock().await;
let mut batches = vec![];
while queue
.iter()
.map(|(_, content)| content.len())
.sum::<usize>()
> max_retained_bytes
{
queue.sort_unstable_by_key(|(_, content)| content.len());
let mut batch = vec![];
let mut batch_size = 0;
for (uploader_chunk_id, content) in std::mem::take(&mut *queue).into_iter().rev() {
if content.len() <= MAX_CHUNK_SIZE - batch_size {
batch_size += content.len();
batch.push((uploader_chunk_id, content));
} else {
queue.push((uploader_chunk_id, content));
}
}
batches.push(batch);
}
try_join_all(batches.into_iter().map(|chunks| async move {
let (uploader_chunk_ids, chunk_contents): (Vec<_>, Vec<_>) = chunks.into_iter().unzip();
let chunk_sizes: Vec<usize> = chunk_contents.iter().map(|c| c.len()).collect();
let canister_chunk_ids = create_chunks(
&self.canister,
&self.batch_id,
chunk_contents,
semaphores,
progress,
)
.await?;
let mut id_map = self.id_mapping.lock().await;
let mut canister_sizes = self.canister_chunk_sizes.lock().await;
for ((uploader_id, canister_id), chunk_size) in uploader_chunk_ids
.into_iter()
.zip(canister_chunk_ids.into_iter())
.zip(chunk_sizes.into_iter())
{
id_map.insert(uploader_id, canister_id.clone());
canister_sizes.insert(canister_id, chunk_size);
}
Ok(())
}))
.await?;
Ok(())
}
}
#[allow(clippy::too_many_arguments)]
async fn make_project_asset_encoding(
chunk_upload_target: Option<&ChunkUploader<'_>>,
asset_descriptor: &AssetDescriptor,
canister_assets: &HashMap<String, AssetDetails>,
content: &Content,
content_encoding: &str,
semaphores: &Semaphores,
logger: &Logger,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<ProjectAssetEncoding, CreateChunkError> {
let sha256 = content.sha256();
let already_in_place = if let Some(canister_asset) = canister_assets.get(&asset_descriptor.key)
{
if canister_asset.content_type != content.media_type.to_string() {
false
} else if let Some(canister_asset_encoding_sha256) = canister_asset
.encodings
.iter()
.find(|details| details.content_encoding == content_encoding)
.and_then(|details| details.sha256.as_ref())
{
canister_asset_encoding_sha256 == &sha256
} else {
false
}
} else {
false
};
let uploader_chunk_ids = if already_in_place {
debug!(
logger,
" {}{} ({} bytes) sha {} is already installed",
&asset_descriptor.key,
content_encoding_descriptive_suffix(content_encoding),
content.data.len(),
hex::encode(&sha256),
);
vec![]
} else if let Some(target) = chunk_upload_target {
upload_content_chunks(
target,
asset_descriptor,
content,
&sha256,
content_encoding,
semaphores,
logger,
progress,
)
.await?
} else {
debug!(
logger,
" {}{} ({} bytes) sha {} will be uploaded",
&asset_descriptor.key,
content_encoding_descriptive_suffix(content_encoding),
content.data.len(),
hex::encode(&sha256),
);
vec![]
};
Ok(ProjectAssetEncoding {
uploader_chunk_ids,
sha256,
already_in_place,
})
}
#[allow(clippy::too_many_arguments)]
async fn make_encoding(
chunk_upload_target: Option<&ChunkUploader<'_>>,
asset_descriptor: &AssetDescriptor,
canister_assets: &HashMap<String, AssetDetails>,
content: &Content,
encoder: &ContentEncoder,
force_encoding: bool,
semaphores: &Semaphores,
logger: &Logger,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<Option<(String, ProjectAssetEncoding)>, CreateEncodingError> {
match encoder {
ContentEncoder::Identity => {
let identity_asset_encoding = make_project_asset_encoding(
chunk_upload_target,
asset_descriptor,
canister_assets,
content,
CONTENT_ENCODING_IDENTITY,
semaphores,
logger,
progress,
)
.await
.map_err(CreateEncodingError::CreateChunkFailed)?;
Ok(Some((
CONTENT_ENCODING_IDENTITY.to_string(),
identity_asset_encoding,
)))
}
encoder => {
let encoded = content.encode(encoder).map_err(|e| {
EncodeContentFailed(asset_descriptor.key.clone(), encoder.to_owned(), e)
})?;
if force_encoding || encoded.data.len() < content.data.len() {
let content_encoding = format!("{encoder}");
let project_asset_encoding = make_project_asset_encoding(
chunk_upload_target,
asset_descriptor,
canister_assets,
&encoded,
&content_encoding,
semaphores,
logger,
progress,
)
.await
.map_err(CreateEncodingError::CreateChunkFailed)?;
Ok(Some((content_encoding, project_asset_encoding)))
} else {
Ok(None)
}
}
}
}
async fn make_encodings(
chunk_upload_target: Option<&ChunkUploader<'_>>,
asset_descriptor: &AssetDescriptor,
canister_assets: &HashMap<String, AssetDetails>,
content: &Content,
semaphores: &Semaphores,
logger: &Logger,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<HashMap<String, ProjectAssetEncoding>, CreateEncodingError> {
let encoders = asset_descriptor
.config
.encodings
.clone()
.unwrap_or_else(|| default_encoders(&content.media_type));
let force_encoding = !encoders.contains(&ContentEncoder::Identity);
let encoding_futures: Vec<_> = encoders
.iter()
.map(|encoder| {
make_encoding(
chunk_upload_target,
asset_descriptor,
canister_assets,
content,
encoder,
force_encoding,
semaphores,
logger,
progress,
)
})
.collect();
let encodings = try_join_all(encoding_futures).await?;
let mut result: HashMap<String, ProjectAssetEncoding> = HashMap::new();
for (key, value) in encodings.into_iter().flatten() {
result.insert(key, value);
}
Ok(result)
}
async fn make_project_asset(
chunk_upload_target: Option<&ChunkUploader<'_>>,
asset_descriptor: AssetDescriptor,
canister_assets: &HashMap<String, AssetDetails>,
semaphores: &Semaphores,
logger: &Logger,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<ProjectAsset, CreateProjectAssetError> {
let file_size = crate::fs::metadata(&asset_descriptor.source)?.len();
let permits = (file_size.div_ceil(1000000) as usize).clamp(1, MAX_COST_SINGLE_FILE_MB);
let _releaser = semaphores.file.acquire(permits).await;
let content = Content::load(&asset_descriptor.source)
.map_err(CreateProjectAssetError::LoadContentFailed)?;
let encodings = make_encodings(
chunk_upload_target,
&asset_descriptor,
canister_assets,
&content,
semaphores,
logger,
progress,
)
.await?;
if let Some(progress) = progress {
progress.increment_complete_assets();
}
Ok(ProjectAsset {
asset_descriptor,
media_type: content.media_type,
encodings,
})
}
pub(crate) async fn make_project_assets(
chunk_upload_target: Option<&ChunkUploader<'_>>,
asset_descriptors: Vec<AssetDescriptor>,
canister_assets: &HashMap<String, AssetDetails>,
mode: Mode,
logger: &Logger,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<HashMap<String, ProjectAsset>, CreateProjectAssetError> {
let semaphores = Semaphores::new();
if let Some(progress) = progress {
progress.set_total_assets(asset_descriptors.len());
}
let project_asset_futures: Vec<_> = asset_descriptors
.iter()
.map(|loc| {
make_project_asset(
chunk_upload_target,
loc.clone(),
canister_assets,
&semaphores,
logger,
progress,
)
})
.collect();
let project_assets = try_join_all(project_asset_futures).await?;
if let Some(uploader) = chunk_upload_target {
uploader
.finalize_upload(&semaphores, mode, progress)
.await
.map_err(|err| {
CreateProjectAssetError::CreateEncodingError(
CreateEncodingError::CreateChunkFailed(err),
)
})?;
}
let mut hm = HashMap::new();
for project_asset in project_assets {
hm.insert(project_asset.asset_descriptor.key.clone(), project_asset);
}
Ok(hm)
}
async fn upload_content_chunks(
chunk_uploader: &ChunkUploader<'_>,
asset_descriptor: &AssetDescriptor,
content: &Content,
sha256: &Vec<u8>,
content_encoding: &str,
semaphores: &Semaphores,
logger: &Logger,
progress: Option<&dyn AssetSyncProgressRenderer>,
) -> Result<Vec<usize>, CreateChunkError> {
if content.data.is_empty() {
let empty = vec![];
let chunk_id = chunk_uploader
.create_chunk(&empty, semaphores, progress)
.await?;
debug!(
logger,
" {}{} 1/1 (0 bytes) sha {}",
&asset_descriptor.key,
content_encoding_descriptive_suffix(content_encoding),
hex::encode(sha256)
);
return Ok(vec![chunk_id]);
}
if let Some(progress) = progress {
progress.add_total_bytes(content.data.len());
}
let count = content.data.len().div_ceil(MAX_CHUNK_SIZE);
let chunks_futures: Vec<_> = content
.data
.chunks(MAX_CHUNK_SIZE)
.enumerate()
.map(|(i, data_chunk)| {
chunk_uploader
.create_chunk(data_chunk, semaphores, progress)
.map_ok(move |chunk_id| {
debug!(
logger,
" {}{} {}/{} ({} bytes) sha {} {}",
&asset_descriptor.key,
content_encoding_descriptive_suffix(content_encoding),
i + 1,
count,
data_chunk.len(),
hex::encode(sha256),
&asset_descriptor.config
);
debug!(logger, "{:?}", &asset_descriptor.config);
chunk_id
})
})
.collect();
try_join_all(chunks_futures).await
}
fn content_encoding_descriptive_suffix(content_encoding: &str) -> String {
if content_encoding == CONTENT_ENCODING_IDENTITY {
"".to_string()
} else {
format!(" ({content_encoding})")
}
}
fn default_encoders(media_type: &Mime) -> Vec<ContentEncoder> {
match (media_type.type_(), media_type.subtype()) {
(mime::TEXT, _) | (_, mime::JAVASCRIPT) | (_, mime::HTML) => {
vec![ContentEncoder::Identity, ContentEncoder::Gzip]
}
_ => vec![ContentEncoder::Identity],
}
}