Skip to main content

azure_storage_blob/clients/
blob_client.rs

1// Copyright (c) Microsoft Corporation. All rights reserved.
2// Licensed under the MIT License.
3
4pub 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    /// Creates a new BlobClient from a blob URL.
32    ///
33    /// # Arguments
34    ///
35    /// * `blob_url` - The full URL of the blob, for example `https://myaccount.blob.core.windows.net/mycontainer/myblob`.
36    ///   The caller is responsible for percent-encoding the URL correctly; it will be used as-is.
37    /// * `credential` - An optional implementation of [`TokenCredential`] that can provide an Entra ID token to use when authenticating.
38    /// * `options` - Optional configuration for the client.
39    #[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        // Storage endpoints must be base URLs.
46        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    /// Returns a new instance of AppendBlobClient.
87    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    /// Returns a new instance of BlockBlobClient.
97    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    /// Returns a new instance of PageBlobClient.
107    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    /// Gets the URL of the resource this client is configured for.
117    pub fn url(&self) -> &Url {
118        &self.endpoint
119    }
120
121    /// Creates a new BlobClient targeting a specific blob version.
122    ///
123    /// # Arguments
124    ///
125    /// * `version_id` - The version ID of the blob to target.
126    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    /// Creates a new BlobClient targeting a specific blob snapshot.
143    ///
144    /// # Arguments
145    ///
146    /// * `snapshot` - The snapshot ID of the blob to target.
147    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    /// Downloads a blob and its contents from the service.
164    ///
165    /// This operation performs a managed (multi-part) download, splitting the blob into
166    /// parallel range requests for better performance on large blobs. The returned
167    /// [`BlobClientDownloadResult::body`] contains the complete blob data, while
168    /// [`BlobClientDownloadResult::properties`] and [`BlobClientDownloadResult::headers`]
169    /// reflect only the initial response's metadata and properties.
170    ///
171    /// If the streamed bytes of this blob are to be collected into contiguous memory,
172    /// consider instead calling [`BlobClient::download_into`] with a pre-allocated buffer
173    /// to avoid unnecessary copies and allocations.
174    ///
175    /// # Arguments
176    ///
177    /// * `options` - Optional configuration for the request.
178    ///
179    /// # Notes
180    ///
181    /// By default, storage clients create their HTTP transport via
182    /// [`azure_core::http::new_http_client()`] with automatic decompression disabled.
183    /// If you set a custom transport in [`BlobClientOptions`] without also disabling
184    /// automatic decompression, partitioned downloads may not succeed.
185    #[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    /// Downloads a blob and its contents from the service.
212    ///
213    /// This operation performs a managed (multi-part) download, splitting the blob into
214    /// parallel range requests for better performance on large blobs. The downloaded bytes are
215    /// written directly into the provided `buffer`.
216    ///
217    /// # Arguments
218    ///
219    /// * `buffer` - Destination buffer to write the downloaded blob data into.
220    /// * `options` - Optional configuration for the request.
221    ///
222    /// # Notes
223    ///
224    /// By default, storage clients create their HTTP transport via
225    /// [`azure_core::http::new_http_client()`] with automatic decompression disabled.
226    /// If you set a custom transport in [`BlobClientOptions`] without also disabling
227    /// automatic decompression, partitioned downloads may not succeed.
228    #[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    /// Uploads content to a block blob, overwriting any existing blob by default.
265    ///
266    /// Updating an existing block blob overwrites any existing metadata on the blob. Use [`BlobClientUploadOptions::if_not_exists()`] to fail instead of overwriting.
267    /// To perform a partial update of the content of a block blob, use [`BlockBlobClient::stage_block()`] and [`BlockBlobClient::commit_block_list()`] directly.
268    ///
269    /// # Arguments
270    ///
271    /// * `content` - The content to upload.
272    /// * `options` - Optional parameters for the request.
273    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    /// Checks if the blob exists.
282    ///
283    /// Returns `true` if the blob exists, `false` if the blob does not exist, and propagates all other errors.
284    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                // Propagate all other error types.
297                _ => 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}