azure_storage_blob/clients/
blob_client.rs1pub use crate::generated::clients::{BlobClient, BlobClientOptions};
5
6use crate::{
7 generated::{
8 clients::BlobClient as GeneratedBlobClient, models::BlobClientDownloadInternalOptions,
9 },
10 models::{
11 BlobClientDownloadIntoResult, BlobClientDownloadOptions, BlobClientDownloadResult,
12 BlobClientUploadOptions, BlobClientUploadResult, BlobDownloadProperties, HttpRange,
13 StorageErrorCode,
14 },
15 partitioned_transfer::{self, PartitionedDownloadBehavior},
16 AppendBlobClient, BlockBlobClient, PageBlobClient,
17};
18use async_trait::async_trait;
19use azure_core::{
20 credentials::TokenCredential,
21 error::ErrorKind,
22 http::{
23 policies::{auth::BearerTokenAuthorizationPolicy, Policy},
24 AsyncRawResponse, Etag, NoFormat, Pipeline, RequestContent, StatusCode, Url, UrlExt,
25 },
26 tracing, Bytes, Result,
27};
28use std::{ops::Range, sync::Arc};
29
30impl BlobClient {
31 #[tracing::new("Storage.Blob.Blob")]
40 pub fn new(
41 blob_url: Url,
42 credential: Option<Arc<dyn TokenCredential>>,
43 options: Option<BlobClientOptions>,
44 ) -> Result<Self> {
45 if blob_url.cannot_be_a_base() {
47 return Err(azure_core::Error::with_message(
48 azure_core::error::ErrorKind::Other,
49 format!("{blob_url} is not a valid base URL"),
50 ));
51 }
52
53 let mut options = options.unwrap_or_default();
54 super::apply_client_defaults(&mut options.client_options);
55
56 let mut per_retry_policies: Vec<Arc<dyn Policy>> = Vec::default();
57 if let Some(token_credential) = credential {
58 if !blob_url.scheme().starts_with("https") {
59 return Err(azure_core::Error::with_message(
60 azure_core::error::ErrorKind::Other,
61 format!("{blob_url} must use https"),
62 ));
63 }
64 per_retry_policies.push(Arc::new(BearerTokenAuthorizationPolicy::new(
65 token_credential,
66 vec!["https://storage.azure.com/.default"],
67 )));
68 }
69
70 let pipeline = Pipeline::new(
71 option_env!("CARGO_PKG_NAME"),
72 option_env!("CARGO_PKG_VERSION"),
73 options.client_options.clone(),
74 Vec::default(),
75 per_retry_policies,
76 None,
77 );
78
79 Ok(Self {
80 endpoint: blob_url,
81 version: options.version,
82 pipeline,
83 })
84 }
85
86 pub fn append_blob_client(&self) -> AppendBlobClient {
88 AppendBlobClient {
89 endpoint: self.endpoint.clone(),
90 pipeline: self.pipeline.clone(),
91 version: self.version.clone(),
92 tracer: self.tracer.clone(),
93 }
94 }
95
96 pub fn block_blob_client(&self) -> BlockBlobClient {
98 BlockBlobClient {
99 endpoint: self.endpoint.clone(),
100 pipeline: self.pipeline.clone(),
101 version: self.version.clone(),
102 tracer: self.tracer.clone(),
103 }
104 }
105
106 pub fn page_blob_client(&self) -> PageBlobClient {
108 PageBlobClient {
109 endpoint: self.endpoint.clone(),
110 pipeline: self.pipeline.clone(),
111 version: self.version.clone(),
112 tracer: self.tracer.clone(),
113 }
114 }
115
116 pub fn url(&self) -> &Url {
118 &self.endpoint
119 }
120
121 pub fn with_version(&self, version_id: &str) -> Result<Self> {
127 let mut versioned_endpoint = self.endpoint.clone();
128 {
129 let mut query_builder = versioned_endpoint.query_builder();
130 query_builder.set_pair("versionid", version_id);
131 query_builder.build();
132 }
133
134 Ok(Self {
135 endpoint: versioned_endpoint,
136 pipeline: self.pipeline.clone(),
137 version: self.version.clone(),
138 tracer: self.tracer.clone(),
139 })
140 }
141
142 pub fn with_snapshot(&self, snapshot: &str) -> Result<Self> {
148 let mut snapshot_endpoint = self.endpoint.clone();
149 {
150 let mut query_builder = snapshot_endpoint.query_builder();
151 query_builder.set_pair("snapshot", snapshot);
152 query_builder.build();
153 }
154
155 Ok(Self {
156 endpoint: snapshot_endpoint,
157 pipeline: self.pipeline.clone(),
158 version: self.version.clone(),
159 tracer: self.tracer.clone(),
160 })
161 }
162
163 #[tracing::function("Storage.Blob.Blob.download")]
186 pub async fn download(
187 &self,
188 options: Option<BlobClientDownloadOptions<'_>>,
189 ) -> Result<BlobClientDownloadResult> {
190 let options = options.unwrap_or_default();
191 let parallel = options
192 .parallel
193 .unwrap_or_else(crate::partitioned_transfer::defaults::default_concurrency);
194 let partition_size = options
195 .partition_size
196 .unwrap_or(crate::partitioned_transfer::defaults::DEFAULT_DOWNLOAD_PARTITION_SIZE);
197 let range = options.range.clone();
198 let inner_client = GeneratedBlobClient {
199 endpoint: self.endpoint.clone(),
200 pipeline: self.pipeline.clone(),
201 version: self.version.clone(),
202 tracer: self.tracer.clone(),
203 };
204 let behavior = BlobClientDownloadBehavior::new(inner_client, options.into());
205 let response =
206 partitioned_transfer::download(range, parallel, partition_size, Arc::new(behavior))
207 .await?;
208 BlobClientDownloadResult::from_headers(response)
209 }
210
211 #[tracing::function("Storage.Blob.Blob.download_into")]
229 pub async fn download_into(
230 &self,
231 buffer: &mut [u8],
232 options: Option<BlobClientDownloadOptions<'_>>,
233 ) -> Result<BlobClientDownloadIntoResult> {
234 let options = options.unwrap_or_default();
235 let parallel = options
236 .parallel
237 .unwrap_or_else(crate::partitioned_transfer::defaults::default_concurrency);
238 let partition_size = options
239 .partition_size
240 .unwrap_or(crate::partitioned_transfer::defaults::DEFAULT_DOWNLOAD_PARTITION_SIZE);
241 let range = options.range.clone();
242 let inner_client = GeneratedBlobClient {
243 endpoint: self.endpoint.clone(),
244 pipeline: self.pipeline.clone(),
245 version: self.version.clone(),
246 tracer: self.tracer.clone(),
247 };
248 let behavior = BlobClientDownloadBehavior::new(inner_client, options.into());
249 let (_, headers, len) = partitioned_transfer::download_into(
250 buffer,
251 range,
252 parallel,
253 partition_size,
254 Arc::new(behavior),
255 )
256 .await?;
257 Ok(BlobClientDownloadIntoResult {
258 len,
259 properties: BlobDownloadProperties::from_headers(&headers)?,
260 headers,
261 })
262 }
263
264 pub async fn upload(
274 &self,
275 content: RequestContent<Bytes, NoFormat>,
276 options: Option<BlobClientUploadOptions<'_>>,
277 ) -> Result<BlobClientUploadResult> {
278 self.block_blob_client().upload(content, options).await
279 }
280
281 pub async fn exists(&self) -> Result<bool> {
285 match self.get_properties(None).await {
286 Ok(_) => Ok(true),
287 Err(e) if e.http_status() == Some(StatusCode::NotFound) => match e.kind() {
288 ErrorKind::HttpResponse {
289 error_code: Some(error_code),
290 ..
291 } if error_code == StorageErrorCode::BlobNotFound.as_ref()
292 || error_code == StorageErrorCode::ContainerNotFound.as_ref() =>
293 {
294 Ok(false)
295 }
296 _ => Err(e),
298 },
299 Err(e) => Err(e),
300 }
301 }
302}
303
304struct BlobClientDownloadBehavior<'a> {
305 client: GeneratedBlobClient,
306 options: BlobClientDownloadInternalOptions<'a>,
307}
308
309impl<'a> BlobClientDownloadBehavior<'a> {
310 fn new(client: GeneratedBlobClient, options: BlobClientDownloadInternalOptions<'a>) -> Self {
311 Self { client, options }
312 }
313}
314
315#[async_trait]
316impl PartitionedDownloadBehavior for BlobClientDownloadBehavior<'_> {
317 async fn transfer_range(
318 &self,
319 range: Option<Range<usize>>,
320 etag_lock: Option<Etag>,
321 ) -> Result<AsyncRawResponse> {
322 let mut opt = self.options.clone();
323 opt.range = range.map(HttpRange::from);
324 if let Some(etag) = etag_lock {
325 opt.if_match = Some(etag);
326 opt.if_none_match = None;
327 opt.if_modified_since = None;
328 opt.if_unmodified_since = None;
329 opt.if_tags = None;
330 }
331 self.client
332 .download_internal(Some(opt))
333 .await
334 .map(AsyncRawResponse::from)
335 }
336}