1pub use crate::generated::clients::{BlockBlobClient, BlockBlobClientOptions};
5
6use crate::{
7 generated::models::{
8 BlockBlobClientCommitBlockListResultHeaders, BlockBlobClientUploadInternalOptions,
9 BlockBlobClientUploadInternalResultHeaders,
10 },
11 models::{
12 BlockBlobClientCommitBlockListOptions, BlockBlobClientStageBlockOptions,
13 BlockBlobClientUploadOptions, BlockBlobClientUploadResult, BlockLookupList,
14 },
15 partitioned_transfer::{self, PartitionedUploadBehavior},
16};
17use async_trait::async_trait;
18use azure_core::{
19 credentials::TokenCredential,
20 http::{
21 policies::{auth::BearerTokenAuthorizationPolicy, Policy},
22 Body, NoFormat, Pipeline, RequestContent, Url,
23 },
24 tracing, Bytes, Result, Uuid,
25};
26use futures::lock::Mutex;
27use std::sync::Arc;
28
29impl BlockBlobClient {
30 #[tracing::new("Storage.Blob.BlockBlob")]
39 pub fn new(
40 blob_url: Url,
41 credential: Option<Arc<dyn TokenCredential>>,
42 options: Option<BlockBlobClientOptions>,
43 ) -> Result<Self> {
44 if blob_url.cannot_be_a_base() {
46 return Err(azure_core::Error::with_message(
47 azure_core::error::ErrorKind::Other,
48 format!("{blob_url} is not a valid base URL"),
49 ));
50 }
51
52 let mut options = options.unwrap_or_default();
53 super::apply_client_defaults(&mut options.client_options);
54
55 let mut per_retry_policies: Vec<Arc<dyn Policy>> = Vec::default();
56 if let Some(token_credential) = credential {
57 if !blob_url.scheme().starts_with("https") {
58 return Err(azure_core::Error::with_message(
59 azure_core::error::ErrorKind::Other,
60 format!("{blob_url} must use https"),
61 ));
62 }
63 per_retry_policies.push(Arc::new(BearerTokenAuthorizationPolicy::new(
64 token_credential,
65 vec!["https://storage.azure.com/.default"],
66 )));
67 }
68
69 let pipeline = Pipeline::new(
70 option_env!("CARGO_PKG_NAME"),
71 option_env!("CARGO_PKG_VERSION"),
72 options.client_options.clone(),
73 Vec::default(),
74 per_retry_policies,
75 None,
76 );
77
78 Ok(Self {
79 endpoint: blob_url,
80 version: options.version,
81 pipeline,
82 })
83 }
84
85 pub fn url(&self) -> &Url {
87 &self.endpoint
88 }
89
90 #[tracing::function("Storage.Blob.BlockBlob.upload")]
100 pub async fn upload(
101 &self,
102 content: RequestContent<Bytes, NoFormat>,
103 options: Option<BlockBlobClientUploadOptions<'_>>,
104 ) -> Result<BlockBlobClientUploadResult> {
105 let options = options.unwrap_or_default();
106 let parallel = options
107 .parallel
108 .unwrap_or_else(crate::partitioned_transfer::defaults::default_concurrency);
109 let partition_size = options
110 .partition_size
111 .unwrap_or(crate::partitioned_transfer::defaults::DEFAULT_UPLOAD_PARTITION_SIZE);
112 let oneshot_options = BlockBlobClientUploadInternalOptions {
114 blob_cache_control: options.blob_cache_control.clone(),
115 blob_content_disposition: options.blob_content_disposition.clone(),
116 blob_content_encoding: options.blob_content_encoding.clone(),
117 blob_content_language: options.blob_content_language.clone(),
118 blob_content_md5: options.blob_content_md5.clone(),
119 blob_content_type: options.blob_content_type.clone(),
120 blob_tags_string: options.blob_tags_string.clone(),
121 encryption_algorithm: options.encryption_algorithm,
122 encryption_key: options.encryption_key.clone(),
123 encryption_key_sha256: options.encryption_key_sha256.clone(),
124 encryption_scope: options.encryption_scope.clone(),
125 if_match: options.if_match.clone(),
126 if_modified_since: options.if_modified_since,
127 if_none_match: options.if_none_match.clone(),
128 if_tags: options.if_tags.clone(),
129 if_unmodified_since: options.if_unmodified_since,
130 immutability_policy_expiry: options.immutability_policy_expiry,
131 immutability_policy_mode: options.immutability_policy_mode,
132 lease_id: options.lease_id.clone(),
133 legal_hold: options.legal_hold,
134 metadata: options.metadata.clone(),
135 method_options: options.method_options.clone(),
136 structured_body_type: None,
137 structured_content_length: None,
138 tier: options.tier.clone(),
139 timeout: options.per_request_timeout,
140 transactional_content_crc64: None,
141 transactional_content_md5: None,
142 };
143 let stage_block_options = BlockBlobClientStageBlockOptions {
144 encryption_algorithm: options.encryption_algorithm,
145 encryption_key: options.encryption_key.clone(),
146 encryption_key_sha256: options.encryption_key_sha256.clone(),
147 encryption_scope: options.encryption_scope.clone(),
148 lease_id: options.lease_id.clone(),
149 method_options: options.method_options.clone(),
150 timeout: options.per_request_timeout,
151 transactional_content_crc64: None,
152 transactional_content_md5: None,
153 };
154 let commit_block_list_options = BlockBlobClientCommitBlockListOptions {
155 blob_cache_control: options.blob_cache_control,
156 blob_content_disposition: options.blob_content_disposition,
157 blob_content_encoding: options.blob_content_encoding,
158 blob_content_language: options.blob_content_language,
159 blob_content_md5: options.blob_content_md5,
160 blob_content_type: options.blob_content_type,
161 blob_tags_string: options.blob_tags_string,
162 encryption_algorithm: options.encryption_algorithm,
163 encryption_key: options.encryption_key,
164 encryption_key_sha256: options.encryption_key_sha256,
165 encryption_scope: options.encryption_scope,
166 if_match: options.if_match,
167 if_modified_since: options.if_modified_since,
168 if_none_match: options.if_none_match,
169 if_tags: options.if_tags,
170 if_unmodified_since: options.if_unmodified_since,
171 immutability_policy_expiry: options.immutability_policy_expiry,
172 immutability_policy_mode: options.immutability_policy_mode,
173 lease_id: options.lease_id,
174 legal_hold: options.legal_hold,
175 metadata: options.metadata,
176 method_options: options.method_options,
177 tier: options.tier,
178 timeout: options.per_request_timeout,
179 transactional_content_crc64: None,
180 transactional_content_md5: None,
181 };
182 let behavior = BlockBlobClientUploadBehavior::new(
183 self,
184 oneshot_options,
185 stage_block_options,
186 commit_block_list_options,
187 );
188 partitioned_transfer::upload(content.into(), parallel, partition_size, &behavior).await?;
189 behavior.result.into_inner().ok_or_else(|| {
190 azure_core::Error::with_message(
191 azure_core::error::ErrorKind::Other,
192 "Upload completed without setting result.",
193 )
194 })
195 }
196}
197
198struct BlockInfo {
199 offset: u64,
200 block_id: Uuid,
201}
202
203struct BlockBlobClientUploadBehavior<'c, 'opt> {
204 client: &'c BlockBlobClient,
205 oneshot_options: BlockBlobClientUploadInternalOptions<'opt>,
206 stage_block_options: BlockBlobClientStageBlockOptions<'opt>,
207 commit_block_list_options: BlockBlobClientCommitBlockListOptions<'opt>,
208 blocks: Mutex<Vec<BlockInfo>>,
209 result: Mutex<Option<BlockBlobClientUploadResult>>,
210}
211
212impl<'c, 'opt> BlockBlobClientUploadBehavior<'c, 'opt> {
213 fn new(
214 client: &'c BlockBlobClient,
215 oneshot_options: BlockBlobClientUploadInternalOptions<'opt>,
216 stage_block_options: BlockBlobClientStageBlockOptions<'opt>,
217 commit_block_list_options: BlockBlobClientCommitBlockListOptions<'opt>,
218 ) -> Self {
219 Self {
220 client,
221 oneshot_options,
222 stage_block_options,
223 commit_block_list_options,
224 blocks: Mutex::new(vec![]),
225 result: Mutex::new(None),
226 }
227 }
228}
229
230#[async_trait]
231impl PartitionedUploadBehavior for BlockBlobClientUploadBehavior<'_, '_> {
232 async fn transfer_oneshot(&self, content: Body) -> Result<()> {
233 let Some(content_len) = content.len() else {
237 return Err(azure_core::Error::with_message(
238 azure_core::error::ErrorKind::Io,
239 "length unknown",
240 ));
241 };
242 let rsp = self
243 .client
244 .upload_internal(
245 content.into(),
246 content_len,
247 Some(self.oneshot_options.clone()),
248 )
249 .await?;
250 *self.result.lock().await = Some(BlockBlobClientUploadResult {
251 content_md5: rsp.content_md5()?,
252 content_crc64: rsp.content_crc64()?,
253 encryption_key_sha256: rsp.encryption_key_sha256()?,
254 encryption_scope: rsp.encryption_scope()?,
255 etag: rsp.etag()?,
256 is_server_encrypted: rsp.is_server_encrypted()?,
257 last_modified: rsp.last_modified()?,
258 version_id: rsp.version_id()?,
259 raw_response: rsp.to_raw_response(),
260 });
261 Ok(())
262 }
263
264 async fn transfer_partition(&self, offset: u64, content: Body) -> Result<()> {
265 let Some(content_len) = content.len() else {
269 return Err(azure_core::Error::with_message(
270 azure_core::error::ErrorKind::Io,
271 "length unknown",
272 ));
273 };
274 let block_id = Uuid::new_v4();
275 {
276 self.blocks
277 .lock()
278 .await
279 .push(BlockInfo { offset, block_id });
280 }
281 self.client
282 .stage_block(
283 block_id.as_bytes(),
284 content_len,
285 content.into(),
286 Some(self.stage_block_options.clone()),
287 )
288 .await?;
289 Ok(())
290 }
291
292 async fn initialize(&self, _content_len: std::option::Option<u64>) -> Result<()> {
293 Ok(())
294 }
295
296 async fn finalize(&self) -> Result<()> {
297 let mut blocks = self.blocks.lock().await;
298 blocks.sort_by_key(|left| left.offset);
299 let blocklist = BlockLookupList {
300 latest: Some(
301 blocks
302 .iter()
303 .map(|bi| bi.block_id.as_bytes().to_vec())
304 .collect(),
305 ),
306 ..Default::default()
307 };
308 let rsp = self
309 .client
310 .commit_block_list(
311 blocklist.try_into()?,
312 Some(self.commit_block_list_options.clone()),
313 )
314 .await?;
315 *self.result.lock().await = Some(BlockBlobClientUploadResult {
316 content_md5: rsp.content_md5()?,
317 content_crc64: rsp.content_crc64()?,
318 encryption_key_sha256: rsp.encryption_key_sha256()?,
319 encryption_scope: rsp.encryption_scope()?,
320 etag: rsp.etag()?,
321 is_server_encrypted: rsp.is_server_encrypted()?,
322 last_modified: rsp.last_modified()?,
323 version_id: rsp.version_id()?,
324 raw_response: rsp.to_raw_response(),
325 });
326 Ok(())
327 }
328}