Skip to main content

objectstore_service/backend/
s3_compatible.rs

1//! S3-compatible backend with generic protocol support.
2
3use std::convert::Infallible;
4use std::error::Error as StdError;
5use std::sync::Arc;
6use std::sync::atomic::Ordering;
7use std::{fmt, io};
8
9use futures_util::{StreamExt, TryStreamExt};
10use objectstore_types::metadata::{HEADER_SIZE, Metadata};
11use objectstore_types::range::{ByteRange, ContentRange};
12use objectstore_types::time::Timestamp;
13use reqwest::header::{HeaderMap, HeaderName, HeaderValue};
14use reqwest::{Body, IntoUrl, Method, RequestBuilder, Response, StatusCode};
15
16use super::extensions::{ResponseExt, SendTraced};
17use crate::backend::common::{
18    self, Backend, DeleteResponse, ExpiryUpdate, GetResponse, MetadataResponse, PutResponse,
19    SetExpiryResponse,
20};
21use crate::backend::extensions::ReqwestResultExt;
22use crate::change_stream::{
23    ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream,
24};
25use crate::error::{Error, ErrorKind, Result, ResultExt as _};
26use crate::id::ObjectId;
27use crate::stream::{ClientStream, counting_stream};
28
29/// Configuration for [`S3CompatibleBackend`].
30///
31/// Supports [Amazon S3] and other S3-compatible services. Authentication is handled via
32/// environment variables (`AWS_ACCESS_KEY_ID` and `AWS_SECRET_ACCESS_KEY`) or IAM roles.
33///
34/// [Amazon S3]: https://aws.amazon.com/s3/
35///
36/// # Example
37///
38/// ```yaml
39/// storage:
40///   type: s3compatible
41///   endpoint: https://s3.amazonaws.com
42///   bucket: my-bucket
43/// ```
44#[derive(Debug, Clone, serde::Deserialize, serde::Serialize)]
45pub struct S3CompatibleConfig {
46    /// S3 endpoint URL.
47    ///
48    /// Examples: `https://s3.amazonaws.com`, `http://localhost:9000` (for MinIO)
49    ///
50    /// # Environment Variables
51    ///
52    /// - `OS__STORAGE__TYPE=s3compatible`
53    /// - `OS__STORAGE__ENDPOINT=https://s3.amazonaws.com`
54    pub endpoint: String,
55
56    /// S3 bucket name.
57    ///
58    /// The bucket must exist before starting the server.
59    ///
60    /// # Environment Variables
61    ///
62    /// - `OS__STORAGE__BUCKET=my-bucket`
63    pub bucket: String,
64
65    /// Reports what this backend stores, for per-usecase cost attribution.
66    ///
67    /// # Default
68    ///
69    /// `None`, which disables reporting for this backend.
70    ///
71    /// # Environment Variables
72    ///
73    /// - `OS__STORAGE__COGS__SHARED_RESOURCE_ID=s3_objectstore`
74    /// - `OS__STORAGE__COGS__SAMPLE_RATE=1.0` (optional)
75    #[serde(default, skip_serializing_if = "Option::is_none")]
76    pub cogs: Option<CostTrackerStreamConfig>,
77}
78
79/// Prefix used for custom metadata in headers for the GCS backend.
80///
81/// See: <https://cloud.google.com/storage/docs/xml-api/reference-headers#xgoogmeta>
82const GCS_CUSTOM_PREFIX: &str = "x-goog-meta-";
83/// Header used to store the expiration time for GCS using the `daysSinceCustomTime` lifecycle
84/// condition.
85///
86/// See: <https://cloud.google.com/storage/docs/xml-api/reference-headers#xgoogcustomtime>
87const GCS_CUSTOM_TIME: &str = "x-goog-custom-time";
88
89/// An authentication token that can be passed as a bearer credential.
90pub trait Token: Send + Sync {
91    /// Returns the token string.
92    fn as_str(&self) -> &str;
93}
94
95/// Provides authentication tokens for S3-compatible requests.
96pub trait TokenProvider: Send + Sync + 'static {
97    /// Error returned when a token cannot be provided.
98    type Error: StdError + Send + Sync + 'static;
99
100    /// Returns a fresh token, fetching or refreshing it as needed.
101    fn get_token(
102        &self,
103    ) -> impl Future<Output = std::result::Result<impl Token, Self::Error>> + Send;
104}
105
106/// Placeholder [`TokenProvider`] for unauthenticated backends.
107#[derive(Debug)]
108pub struct NoToken;
109
110impl TokenProvider for NoToken {
111    type Error = Infallible;
112
113    #[allow(refining_impl_trait)]
114    async fn get_token(&self) -> std::result::Result<NoToken, Infallible> {
115        unimplemented!()
116    }
117}
118impl Token for NoToken {
119    fn as_str(&self) -> &str {
120        unimplemented!()
121    }
122}
123
124/// S3-compatible storage backend with pluggable authentication.
125pub struct S3CompatibleBackend<T> {
126    client: reqwest::Client,
127
128    endpoint: String,
129    bucket: String,
130
131    token_provider: Option<T>,
132
133    change_stream: Arc<dyn ChangeStream>,
134}
135
136impl<T> S3CompatibleBackend<T> {
137    /// Creates a new S3-compatible backend bound to the given bucket.
138    pub fn new(
139        config: S3CompatibleConfig,
140        token_provider: T,
141        streams: &ChangeStreamFactory,
142    ) -> Self {
143        Self::build(config, Some(token_provider), streams)
144    }
145
146    fn build(
147        config: S3CompatibleConfig,
148        token_provider: Option<T>,
149        streams: &ChangeStreamFactory,
150    ) -> Self {
151        let S3CompatibleConfig {
152            endpoint,
153            bucket,
154            cogs,
155        } = config;
156        Self {
157            client: common::reqwest_client(),
158            endpoint,
159            bucket,
160            token_provider,
161            change_stream: streams.build(cogs.as_ref()),
162        }
163    }
164
165    /// Formats the S3 object URL for the given key.
166    fn object_url(&self, id: &ObjectId) -> String {
167        format!("{}/{}/{}", self.endpoint, self.bucket, id.as_storage_path())
168    }
169}
170
171/// Number of bytes the given headers occupy as stored object metadata.
172fn headers_size(headers: &HeaderMap) -> u64 {
173    headers
174        .iter()
175        .map(|(name, value)| name.as_str().len() as u64 + value.len() as u64)
176        .sum()
177}
178
179/// Wraps [`Metadata::to_headers`] with GCS-specific concerns (tombstone + custom-time).
180fn metadata_to_gcs_headers(
181    metadata: &Metadata,
182    prefix: &str,
183) -> Result<HeaderMap, objectstore_types::metadata::Error> {
184    let mut headers = metadata.to_headers(prefix)?;
185
186    // The size is derived from the native `Content-Length` on every read, so it must not be
187    // persisted: metadata updates rewrite *all* stored metadata, and a stored `x-sn-size` key
188    // is rejected by the GCS JSON backend when it deserializes the object.
189    let size = HeaderName::try_from(format!("{prefix}{HEADER_SIZE}"))?;
190    headers.remove(&size);
191
192    // GCS custom-time for lifecycle expiration
193    if let Some(expires_at) = metadata.time_expires {
194        let expires_at = expires_at.as_rfc3339();
195        headers.append(GCS_CUSTOM_TIME, expires_at.to_string().parse()?);
196    }
197    Ok(headers)
198}
199
200impl<T> S3CompatibleBackend<T>
201where
202    T: TokenProvider,
203{
204    /// Creates a request builder with the appropriate authentication.
205    async fn request(&self, method: Method, url: impl IntoUrl) -> Result<RequestBuilder> {
206        let mut builder = self.client.request(method, url);
207        if let Some(provider) = &self.token_provider {
208            builder = builder.bearer_auth(
209                provider
210                    .get_token()
211                    .await
212                    .context(ErrorKind::BackendFailure, "getting S3 authentication token")?
213                    .as_str(),
214            );
215        }
216        Ok(builder)
217    }
218
219    /// Fetches object metadata using the given HTTP method (GET or HEAD) and
220    /// returns it with the response without modifying the object.
221    async fn request_object(
222        &self,
223        method: Method,
224        id: &ObjectId,
225        access_time: Timestamp,
226        range: Option<ByteRange>,
227    ) -> Result<Option<(Metadata, Option<ContentRange>, Response)>> {
228        let object_url = self.object_url(id);
229
230        let mut builder = self.request(method, &object_url).await?;
231        if let Some(r) = range {
232            builder = builder.header(reqwest::header::RANGE, r.to_header_value());
233        }
234        let response = builder
235            .send_traced()
236            .await
237            .reqwest_context("sending an S3 object request")?;
238
239        if response.status() == StatusCode::NOT_FOUND {
240            objectstore_log::debug!("Object not found");
241            response.drain_body().await;
242            return Ok(None);
243        }
244
245        if response.status() == StatusCode::RANGE_NOT_SATISFIABLE {
246            let raw = response
247                .headers()
248                .get(reqwest::header::CONTENT_RANGE)
249                .and_then(|v| v.to_str().ok());
250            let total = raw.and_then(ContentRange::parse_unsatisfiable_total);
251            let err = match total {
252                Some(total) => ErrorKind::RangeNotSatisfiable { total }.into(),
253                None => Error::new(ErrorKind::BackendFailure, "invalid S3 416 Content-Range"),
254            };
255            response.drain_body().await;
256            return Err(err);
257        }
258
259        let response = response.check_error("getting an S3 object").await?;
260
261        let headers = response.headers();
262        let mut metadata = Metadata::from_headers(headers, GCS_CUSTOM_PREFIX)
263            .context(ErrorKind::CorruptData, "decoding S3 object metadata")?;
264
265        let content_range = if response.status() == StatusCode::PARTIAL_CONTENT {
266            let range = headers
267                .get(reqwest::header::CONTENT_RANGE)
268                .and_then(|v| v.to_str().ok())
269                .and_then(|s| s.parse::<ContentRange>().ok())
270                .ok_or_else(|| {
271                    Error::new(ErrorKind::BackendFailure, "missing S3 206 Content-Range")
272                })?;
273            metadata.size = Some(range.total as usize);
274            Some(range)
275        } else {
276            // NB: Read the header rather than `Response::content_length`, which reports the
277            // length of the decoded body and is therefore always zero for a HEAD response.
278            let size = headers
279                .get(reqwest::header::CONTENT_LENGTH)
280                .and_then(|value| value.to_str().ok())
281                .map(|value| value.parse::<usize>())
282                .transpose()
283                .context(ErrorKind::CorruptData, "decoding S3 Content-Length")?;
284
285            if let Some(size) = size {
286                metadata.size = Some(size);
287            } else {
288                objectstore_log::warn!("S3: 200 response missing Content-Length header");
289            }
290            None
291        };
292
293        // Filter already expired objects but leave them to garbage collection
294        if metadata.is_expired(access_time) {
295            objectstore_log::debug!("Object found but past expiry");
296            response.drain_body().await;
297            return Ok(None);
298        }
299
300        Ok(Some((metadata, content_range, response)))
301    }
302
303    /// Issues a request to update the metadata for the given object.
304    async fn update_metadata(
305        &self,
306        id: &ObjectId,
307        metadata: &Metadata,
308        deadline: Timestamp,
309        etag: &HeaderValue,
310    ) -> Result<SetExpiryResponse> {
311        // NB: Meta updates require CopyObject + REPLACE along with *all* metadata. See
312        // https://docs.aws.amazon.com/AmazonS3/latest/API/API_CopyObject.html
313        let request = self
314            .request(Method::PUT, self.object_url(id))
315            .await?
316            .header(
317                "x-amz-copy-source",
318                format!("/{}/{}", self.bucket, id.as_storage_path()),
319            )
320            .header("x-amz-metadata-directive", "REPLACE")
321            .header("x-amz-copy-source-if-match", etag.clone())
322            .headers(
323                metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX)
324                    .context(ErrorKind::InvalidMetadata, "encoding S3 object metadata")?,
325            );
326
327        let response = request.send_traced().await;
328        let response = response.reqwest_context("updating S3 expiration")?;
329        let outcome = match response.status() {
330            StatusCode::NOT_FOUND => Some(SetExpiryResponse::NotFound),
331            StatusCode::CONFLICT | StatusCode::PRECONDITION_FAILED => {
332                Some(SetExpiryResponse::Rejected)
333            }
334            _ => None,
335        };
336        if let Some(outcome) = outcome {
337            response.drain_body().await;
338            return Ok(outcome);
339        }
340        response
341            .check_error("updating S3 expiration")
342            .await?
343            .drain_body()
344            .await;
345
346        Ok(SetExpiryResponse::Satisfied(deadline))
347    }
348}
349
350impl<T> fmt::Debug for S3CompatibleBackend<T> {
351    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
352        f.debug_struct("S3Compatible")
353            .field("client", &self.client)
354            .field("endpoint", &self.endpoint)
355            .field("bucket", &self.bucket)
356            .finish_non_exhaustive()
357    }
358}
359
360impl S3CompatibleBackend<NoToken> {
361    /// Creates a new S3-compatible backend that sends unauthenticated requests.
362    pub fn without_token(config: S3CompatibleConfig, streams: &ChangeStreamFactory) -> Self {
363        Self::build(config, None, streams)
364    }
365}
366
367#[async_trait::async_trait]
368impl<T: TokenProvider> Backend for S3CompatibleBackend<T> {
369    fn name(&self) -> &'static str {
370        "s3-compatible"
371    }
372
373    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
374    async fn put_object(
375        &self,
376        id: &ObjectId,
377        metadata: &Metadata,
378        stream: ClientStream,
379        _access_time: Timestamp,
380    ) -> Result<PutResponse> {
381        objectstore_log::debug!("Writing to s3_compatible backend");
382        let headers = metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX)
383            .context(ErrorKind::InvalidMetadata, "encoding S3 object metadata")?;
384        let metadata_size = headers_size(&headers);
385
386        // A successful PUT does not report the stored size back, so count what we send.
387        let (payload_size, counted) = counting_stream(stream);
388
389        self.request(Method::PUT, self.object_url(id))
390            .await?
391            .headers(headers)
392            .body(Body::wrap_stream(counted))
393            .send_traced()
394            .await
395            .check_error("uploading an S3 object")
396            .await?
397            .drain_body()
398            .await;
399
400        self.change_stream.write(
401            id,
402            metadata_size + payload_size.load(Ordering::Relaxed),
403            metadata.time_expires,
404        );
405
406        Ok(())
407    }
408
409    #[tracing::instrument(level = "debug", skip(self))]
410    async fn get_object(
411        &self,
412        id: &ObjectId,
413        access_time: Timestamp,
414        range: Option<ByteRange>,
415    ) -> Result<GetResponse> {
416        objectstore_log::debug!("Reading from s3_compatible backend");
417
418        let Some((metadata, content_range, response)) = self
419            .request_object(Method::GET, id, access_time, range)
420            .await?
421        else {
422            return Ok(None);
423        };
424
425        let stream = response.bytes_stream().map_err(io::Error::other);
426        Ok(Some((metadata, content_range, stream.boxed())))
427    }
428
429    #[tracing::instrument(level = "debug", skip(self))]
430    async fn get_metadata(
431        &self,
432        id: &ObjectId,
433        access_time: Timestamp,
434    ) -> Result<MetadataResponse> {
435        objectstore_log::debug!("Reading metadata from s3_compatible backend");
436        let response = self
437            .request_object(Method::HEAD, id, access_time, None)
438            .await?;
439        Ok(response.map(|(metadata, _, _)| metadata))
440    }
441
442    #[tracing::instrument(level = "debug", skip(self))]
443    async fn set_expiry(
444        &self,
445        id: &ObjectId,
446        target: ExpiryUpdate,
447        access_time: Timestamp,
448    ) -> Result<SetExpiryResponse> {
449        let Some((mut metadata, _, response)) = self
450            .request_object(Method::HEAD, id, access_time, None)
451            .await?
452        else {
453            return Ok(SetExpiryResponse::NotFound);
454        };
455
456        let etag = response.headers().get(reqwest::header::ETAG).cloned();
457        response.drain_body().await;
458
459        let Some(current_expiry) = metadata.time_expires else {
460            return Ok(SetExpiryResponse::Rejected);
461        };
462        let Some(expire_at) = target.resolve(metadata.time_created, access_time)? else {
463            return Ok(SetExpiryResponse::Rejected);
464        };
465        if current_expiry >= expire_at {
466            return Ok(SetExpiryResponse::Satisfied(expire_at)); // already satisfied
467        }
468
469        let etag = etag.ok_or_else(|| {
470            Error::new(ErrorKind::BackendFailure, "S3 HEAD response missing ETag")
471        })?;
472
473        metadata.expiration_policy = common::extended_expiration_policy(
474            metadata.expiration_policy,
475            metadata.time_created,
476            current_expiry,
477            expire_at,
478        )?;
479        metadata.time_expires = Some(expire_at);
480        let outcome = self
481            .update_metadata(id, &metadata, expire_at, &etag)
482            .await?;
483        if matches!(outcome, SetExpiryResponse::Satisfied(_)) {
484            self.change_stream.update(id, Some(expire_at));
485        }
486
487        Ok(outcome)
488    }
489
490    #[tracing::instrument(level = "debug", skip(self))]
491    async fn delete_object(
492        &self,
493        id: &ObjectId,
494        _access_time: Timestamp,
495    ) -> Result<DeleteResponse> {
496        objectstore_log::debug!("Deleting from s3_compatible backend");
497        let response = self
498            .request(Method::DELETE, self.object_url(id))
499            .await?
500            .send_traced()
501            .await
502            .reqwest_context("sending an S3 delete request")?;
503
504        // S3 deletes are idempotent; they return 204 whether a key existed or not. This
505        // branch catches other 404s, like from a missing bucket.
506        if response.status() == StatusCode::NOT_FOUND {
507            response.drain_body().await;
508            return Ok(());
509        }
510
511        response
512            .check_error("deleting an S3 object")
513            .await?
514            .drain_body()
515            .await;
516
517        // If the object didn't exist in the first place, this emits a spurious message
518        // due to S3 returning 204 to DELETEs whether the object existed or not.
519        self.change_stream.delete(id);
520
521        Ok(())
522    }
523
524    async fn join(&self) {
525        flush_change_stream(&self.change_stream).await;
526    }
527}
528
529#[cfg(test)]
530mod tests {
531    use std::collections::BTreeMap;
532    use std::io::{Read, Write};
533    use std::net::{TcpListener, TcpStream};
534    use std::sync::mpsc;
535    use std::thread;
536    use std::time::Duration;
537
538    use anyhow::Result;
539    use objectstore_types::metadata::ExpirationPolicy;
540    use objectstore_types::scope::{Scope, Scopes};
541
542    use super::*;
543    use crate::backend::common::Backend;
544    use crate::id::ObjectContext;
545    use crate::stream;
546
547    // NB: To run these tests, you need to have SeaweedFS running. This is done
548    // automatically in CI.
549    //
550    // Refer to the readme for how to set up SeaweedFS via devservices.
551
552    fn create_test_backend() -> S3CompatibleBackend<NoToken> {
553        S3CompatibleBackend::without_token(
554            S3CompatibleConfig {
555                endpoint: "http://localhost:8089".into(),
556                bucket: "test-bucket".into(),
557                cogs: None,
558            },
559            &ChangeStreamFactory::default(),
560        )
561    }
562
563    fn make_id() -> ObjectId {
564        ObjectId::random(ObjectContext {
565            usecase: "testing".into(),
566            scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
567        })
568    }
569
570    fn read_http_request(connection: &mut TcpStream) -> String {
571        let mut bytes = Vec::new();
572        let mut byte = [0];
573        while !bytes.ends_with(b"\r\n\r\n") {
574            connection.read_exact(&mut byte).unwrap();
575            bytes.push(byte[0]);
576        }
577        String::from_utf8(bytes).unwrap()
578    }
579
580    fn start_copy_server(
581        copy_status: &'static str,
582        metadata: Metadata,
583    ) -> (String, mpsc::Receiver<String>, thread::JoinHandle<()>) {
584        let listener = TcpListener::bind(("127.0.0.1", 0)).unwrap();
585        let endpoint = format!("http://{}", listener.local_addr().unwrap());
586        let (request_tx, request_rx) = mpsc::channel();
587        let server = thread::spawn(move || {
588            let (mut head, _) = listener.accept().unwrap();
589            assert!(read_http_request(&mut head).starts_with("HEAD "));
590            write!(
591                head,
592                "HTTP/1.1 200 OK\r\nContent-Length: 0\r\nETag: \"etag\"\r\nConnection: close\r\n"
593            )
594            .unwrap();
595            for (name, value) in metadata.to_headers(GCS_CUSTOM_PREFIX).unwrap().iter() {
596                write!(head, "{name}: {}\r\n", value.to_str().unwrap()).unwrap();
597            }
598            write!(head, "\r\n").unwrap();
599            let (mut copy, _) = listener.accept().unwrap();
600            request_tx.send(read_http_request(&mut copy)).unwrap();
601            write!(
602                copy,
603                "HTTP/1.1 {copy_status}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
604            )
605            .unwrap();
606        });
607        (endpoint, request_rx, server)
608    }
609
610    #[tokio::test]
611    async fn update_metadata_uses_conditional_s3_copy() {
612        let created = Timestamp::now();
613        let deadline = created + Duration::from_hours(2);
614        for (status, expected) in [
615            ("200 OK", SetExpiryResponse::Satisfied(deadline)),
616            ("404 Not Found", SetExpiryResponse::NotFound),
617            ("409 Conflict", SetExpiryResponse::Rejected),
618            ("412 Precondition Failed", SetExpiryResponse::Rejected),
619        ] {
620            let (endpoint, request_rx, server) = start_copy_server(
621                status,
622                Metadata {
623                    expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
624                    time_created: Some(created),
625                    time_expires: Some(created + Duration::from_hours(1)),
626                    custom: [("preserved".into(), "yes".into())].into(),
627                    ..Default::default()
628                },
629            );
630            let backend = S3CompatibleBackend::without_token(
631                S3CompatibleConfig {
632                    endpoint,
633                    bucket: "bucket".into(),
634                    cogs: None,
635                },
636                &ChangeStreamFactory::default(),
637            );
638
639            assert_eq!(
640                backend
641                    .set_expiry(
642                        &make_id(),
643                        common::ExpiryTarget::At(deadline).into(),
644                        created
645                    )
646                    .await
647                    .unwrap(),
648                expected
649            );
650            let request = request_rx.recv().unwrap().to_ascii_lowercase();
651            assert!(request.contains("x-amz-copy-source: /bucket/"));
652            assert!(request.contains("x-amz-metadata-directive: replace"));
653            assert!(request.contains("x-amz-copy-source-if-match: \"etag\""));
654            let expected_headers = Metadata {
655                expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(2)),
656                time_created: Some(created),
657                time_expires: Some(deadline),
658                custom: [("preserved".into(), "yes".into())].into(),
659                ..Default::default()
660            }
661            .to_headers(GCS_CUSTOM_PREFIX)
662            .unwrap();
663            for (name, value) in &expected_headers {
664                assert!(request.contains(
665                    &format!("{name}: {}", value.to_str().unwrap()).to_ascii_lowercase()
666                ));
667            }
668            server.join().unwrap();
669        }
670    }
671
672    #[test]
673    fn metadata_to_gcs_headers_omits_size() {
674        let metadata = Metadata {
675            size: Some(4096),
676            ..Default::default()
677        };
678
679        let headers = metadata_to_gcs_headers(&metadata, GCS_CUSTOM_PREFIX).unwrap();
680
681        // Persisting the size would store a key that the GCS JSON backend rejects on read.
682        assert!(headers.get("x-goog-meta-x-sn-size").is_none());
683    }
684
685    #[test]
686    fn metadata_to_gcs_headers_uses_time_expires() {
687        let expires = Timestamp::now() + Duration::from_hours(1);
688        let metadata = Metadata {
689            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
690            time_expires: Some(expires),
691            ..Default::default()
692        };
693
694        let headers = metadata_to_gcs_headers(&metadata, GCS_CUSTOM_PREFIX).unwrap();
695        let custom_time = headers.get(GCS_CUSTOM_TIME).unwrap().to_str().unwrap();
696        let expected = expires.as_rfc3339().to_string();
697        assert_eq!(custom_time, expected);
698    }
699
700    #[test]
701    fn metadata_to_gcs_headers_escapes_unicode() {
702        let metadata = Metadata {
703            filename: Some("réport-📄.pdf".into()),
704            custom: BTreeMap::from_iter([("release".into(), "vérsion-1.0-🚀".into())]),
705            ..Default::default()
706        };
707
708        let headers = metadata_to_gcs_headers(&metadata, GCS_CUSTOM_PREFIX).unwrap();
709        assert_eq!(
710            headers.get("x-goog-meta-x-sn-filename").unwrap(),
711            "r%C3%A9port-%F0%9F%93%84.pdf",
712        );
713        assert_eq!(
714            headers.get("x-goog-meta-x-snme-release").unwrap(),
715            "v%C3%A9rsion-1.0-%F0%9F%9A%80",
716        );
717
718        // The prefixed headers this backend writes are the ones it reads back.
719        let roundtripped = Metadata::from_headers(&headers, GCS_CUSTOM_PREFIX).unwrap();
720        assert_eq!(roundtripped.filename, metadata.filename);
721        assert_eq!(roundtripped.custom, metadata.custom);
722    }
723
724    #[test]
725    fn headers_size_counts_names_and_values() {
726        let mut headers = HeaderMap::new();
727        headers.insert("x-goog-meta-a", "1".parse().unwrap());
728        headers.insert("x-goog-meta-bb", "22".parse().unwrap());
729
730        assert_eq!(
731            headers_size(&headers),
732            ("x-goog-meta-a".len() + 1 + "x-goog-meta-bb".len() + 2) as u64
733        );
734    }
735
736    #[tokio::test]
737    async fn test_get_metadata_nonexistent() -> Result<()> {
738        let backend = create_test_backend();
739        let id = make_id();
740        let result = backend.get_metadata(&id, Timestamp::now()).await?;
741        assert!(result.is_none());
742        Ok(())
743    }
744
745    #[tokio::test]
746    async fn test_get_metadata_reports_size() -> Result<()> {
747        let backend = create_test_backend();
748        let id = make_id();
749        let payload = "hello, world";
750
751        backend
752            .put_object(
753                &id,
754                &Metadata::default(),
755                stream::single(payload),
756                Timestamp::now(),
757            )
758            .await?;
759
760        // The size must come from the `Content-Length` header, not from the (empty) body of
761        // the HEAD response.
762        let metadata = backend
763            .get_metadata(&id, Timestamp::now())
764            .await?
765            .expect("object exists");
766        assert_eq!(metadata.size, Some(payload.len()));
767
768        Ok(())
769    }
770
771    #[tokio::test]
772    #[ignore = "SeaweedFS does not enforce expiration from GCS custom-time metadata"]
773    async fn test_ttl_immediate() -> Result<()> {
774        let backend = create_test_backend();
775        let id = make_id();
776        let metadata = Metadata {
777            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(0)),
778            time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
779            ..Default::default()
780        };
781
782        backend
783            .put_object(
784                &id,
785                &metadata,
786                stream::single("hello, world"),
787                Timestamp::now(),
788            )
789            .await?;
790
791        let get_result = backend.get_object(&id, Timestamp::now(), None).await?;
792        assert!(get_result.is_none());
793
794        let head_result = backend.get_metadata(&id, Timestamp::now()).await?;
795        assert!(head_result.is_none());
796
797        Ok(())
798    }
799
800    #[tokio::test]
801    #[ignore = "SeaweedFS does not enforce expiration from GCS custom-time metadata"]
802    async fn test_tti_immediate() -> Result<()> {
803        let backend = create_test_backend();
804        let id = make_id();
805        let metadata = Metadata {
806            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_secs(0)),
807            time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
808            ..Default::default()
809        };
810
811        backend
812            .put_object(
813                &id,
814                &metadata,
815                stream::single("hello, world"),
816                Timestamp::now(),
817            )
818            .await?;
819
820        let get_result = backend.get_object(&id, Timestamp::now(), None).await?;
821        assert!(get_result.is_none());
822
823        let head_result = backend.get_metadata(&id, Timestamp::now()).await?;
824        assert!(head_result.is_none());
825
826        Ok(())
827    }
828}