objectstore_service/backend/common.rs
1//! Shared trait definition and types for all backends.
2
3use std::fmt;
4use std::num::NonZeroU64;
5use std::time::Duration;
6
7use objectstore_types::metadata::{ExpirationPolicy, Metadata};
8use objectstore_types::range::{ByteRange, ContentRange};
9use objectstore_types::resumable::UploadProgress;
10use objectstore_types::time::Timestamp;
11
12use bytes::Bytes;
13
14use crate::error::{Error, ErrorKind, Result};
15use crate::id::ObjectId;
16use crate::multipart::{
17 AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
18 ListPartsResponse, PartNumber, UploadId, UploadPartResponse,
19};
20use crate::resumable::{BackendToken, Session};
21use crate::stream::{ClientStream, PayloadStream};
22
23/// User agent string used for outgoing requests.
24///
25/// This intentionally has a "sentry" prefix so that it can easily be traced back to us.
26pub const USER_AGENT: &str = concat!("sentry-objectstore/", env!("CARGO_PKG_VERSION"));
27
28/// Backend response for put operations.
29pub type PutResponse = ();
30/// Backend response for get operations.
31pub type GetResponse = Option<(Metadata, Option<ContentRange>, PayloadStream)>;
32/// Backend response for metadata-only get operations.
33pub type MetadataResponse = Option<Metadata>;
34/// Backend response for delete operations.
35pub type DeleteResponse = ();
36
37/// The outcome of an expiry update.
38#[derive(Clone, Copy, Debug, PartialEq, Eq)]
39pub enum SetExpiryResponse {
40 /// The deadline was extended or already satisfied the request.
41 ///
42 /// Contains the resolved requested deadline, not necessarily the stored deadline.
43 Satisfied(Timestamp),
44 /// The object or redirect was observed to be absent or expired.
45 NotFound,
46 /// The update could not be satisfied.
47 ///
48 /// The entry is non-expiring, lacks required creation metadata, or conflicts
49 /// with the conditional update. A failed conditional write does not establish
50 /// absence, even if a concurrent deletion caused it to fail.
51 Rejected,
52}
53
54/// The requested minimum deadline for an expiry update.
55///
56/// [`ExpiryTarget::At`] is already resolved. [`ExpiryTarget::FromCreation`]
57/// is resolved by the backend from the creation time read by the update
58/// operation itself. Zero durations and resolved deadlines in the past remain
59/// valid minimum-deadline requests.
60#[derive(Clone, Copy, Debug, PartialEq, Eq)]
61pub enum ExpiryTarget {
62 /// An absolute deadline.
63 At(Timestamp),
64 /// A deadline relative to the object's creation time.
65 FromCreation(Duration),
66}
67
68impl ExpiryTarget {
69 /// Resolves this target against an optional creation time.
70 ///
71 /// Returns `None` when a creation-relative target has no creation time.
72 /// Creation-relative deadlines use the timestamp's existing rounding and
73 /// clamp to its maximum value on overflow.
74 pub fn resolve(self, time_created: Option<Timestamp>) -> Option<Timestamp> {
75 match self {
76 Self::At(deadline) => Some(deadline),
77 Self::FromCreation(duration) => {
78 time_created.map(|created| created.saturating_add(duration))
79 }
80 }
81 }
82}
83
84/// An expiry target and optional limit on the remaining lifetime it requests.
85///
86/// The limit is measured from the operation's `access_time`, regardless of the target's
87/// anchor. It does not change the policy kind or shorten existing deadlines.
88/// Repeated updates can therefore keep an object alive indefinitely.
89#[derive(Clone, Copy, Debug, PartialEq, Eq)]
90pub struct ExpiryUpdate {
91 /// The requested minimum deadline.
92 pub target: ExpiryTarget,
93 /// Maximum remaining lifetime. `None` means no limit.
94 pub max: Option<Duration>,
95}
96
97impl From<ExpiryTarget> for ExpiryUpdate {
98 fn from(target: ExpiryTarget) -> Self {
99 Self { target, max: None }
100 }
101}
102
103impl ExpiryUpdate {
104 /// Resolves and validates the requested deadline using the stored object metadata.
105 ///
106 /// Returns `None` if the creation-time anchor is required but unavailable.
107 ///
108 /// Returns [`ErrorKind::InvalidMetadata`] if the requested deadline exceeds
109 /// `access_time + max`, even if the existing deadline already satisfies it.
110 /// Timestamp addition rounds up to seconds and saturates at the supported maximum.
111 pub fn resolve(
112 self,
113 time_created: Option<Timestamp>,
114 access_time: Timestamp,
115 ) -> Result<Option<Timestamp>> {
116 let Some(deadline) = self.target.resolve(time_created) else {
117 return Ok(None);
118 };
119
120 if let Some(max) = self.max
121 && deadline > access_time.saturating_add(max)
122 {
123 return Err(Error::new(
124 ErrorKind::InvalidMetadata,
125 "requested expiry exceeds maximum remaining lifetime",
126 ));
127 }
128 Ok(Some(deadline))
129 }
130}
131
132/// Derives a TTL matching an extended deadline, leaving other policies unchanged.
133///
134/// Uses creation time when available, otherwise infers the original lifetime from
135/// the previous deadline and TTL. The inferred creation time is not persisted.
136/// Rounds to the coarsest whole day, hour, or minute within 1% of the duration,
137/// falling back to whole seconds. The deadline itself must remain exact.
138///
139/// Returns corrupt-data errors for inconsistent creation times or duration overflow.
140pub(super) fn extended_expiration_policy(
141 policy: ExpirationPolicy,
142 time_created: Option<Timestamp>,
143 old_expiry: Timestamp,
144 new_expiry: Timestamp,
145) -> Result<ExpirationPolicy> {
146 let ExpirationPolicy::TimeToLive(ttl) = policy else {
147 return Ok(policy);
148 };
149
150 let duration_opt = match time_created {
151 Some(created) => new_expiry.checked_duration_since(created),
152 None => new_expiry
153 .checked_duration_since(old_expiry)
154 .and_then(|extension| ttl.checked_add(extension)),
155 };
156
157 let duration = duration_opt
158 .ok_or_else(|| Error::new(ErrorKind::CorruptData, "invalid TTL expiry metadata"))?;
159 Ok(ExpirationPolicy::TimeToLive(round_duration(duration)))
160}
161
162/// Rounds a duration to the coarsest whole day, hour, or minute within 1%, otherwise seconds.
163fn round_duration(duration: Duration) -> Duration {
164 let mut seconds = duration.as_secs();
165
166 // Match Timestamp's upward rounding for legacy fractional durations.
167 if duration.subsec_nanos() != 0 {
168 seconds = seconds.saturating_add(1);
169 }
170
171 for unit in [86_400, 3_600, 60] {
172 let rounded = seconds.saturating_add(unit / 2) / unit * unit;
173 if rounded.abs_diff(seconds) <= seconds / 100 {
174 return Duration::from_secs(rounded);
175 }
176 }
177
178 Duration::from_secs(seconds)
179}
180
181/// Trait implemented by all storage backends.
182///
183/// Object operations take `access_time`, the timestamp of the caller's operation.
184/// Use it to decide whether an object or redirect has expired, so all steps of
185/// an operation use the same time, including retries and calls to other backends.
186/// An object is expired when its deadline is strictly earlier than `access_time`.
187/// Writes preserve the creation time and deadline in the supplied metadata.
188#[async_trait::async_trait]
189pub trait Backend: fmt::Debug + Send + Sync + 'static {
190 /// The backend name, used for diagnostics.
191 fn name(&self) -> &'static str;
192
193 /// Returns the upload granularity for sessions opened by this backend, in bytes.
194 ///
195 /// A value of zero means these uploads have no granularity. A positive value means that a
196 /// non-final chunk can persist only a multiple of this value. Implementations must reject a
197 /// non-empty, non-final chunk shorter than one unit with [`ErrorKind::ChunkTooSmall`].
198 fn upload_granularity(&self) -> u64 {
199 0
200 }
201
202 /// Stores an object at the given path with the given metadata.
203 async fn put_object(
204 &self,
205 id: &ObjectId,
206 metadata: &Metadata,
207 stream: ClientStream,
208 access_time: Timestamp,
209 ) -> Result<PutResponse>;
210
211 /// Retrieves (part of) an object at the given path, returning its metadata, a description of
212 /// the part being returned, and the payload.
213 async fn get_object(
214 &self,
215 id: &ObjectId,
216 access_time: Timestamp,
217 range: Option<ByteRange>,
218 ) -> Result<GetResponse>;
219
220 /// Retrieves only the metadata for an object, without the payload.
221 async fn get_metadata(
222 &self,
223 id: &ObjectId,
224 access_time: Timestamp,
225 ) -> Result<MetadataResponse> {
226 Ok(self
227 .get_object(id, access_time, None)
228 .await?
229 .map(|(metadata, _range, _stream)| metadata))
230 }
231
232 /// Extends the deadline of an existing object with expiration policy.
233 ///
234 /// Actual extensions also update TTL durations to approximately match the time since
235 /// creation. TTI durations, payload, and other metadata remain unchanged.
236 ///
237 /// Returns [`SetExpiryResponse::Satisfied`] when extended or already satisfied,
238 /// [`SetExpiryResponse::NotFound`] when observed absent or expired, or
239 /// [`SetExpiryResponse::Rejected`] when ineligible or conflicting.
240 /// Limit violations return [`ErrorKind::InvalidMetadata`]; backend failures also
241 /// return errors. Limits are checked before reporting an already-satisfied request.
242 async fn set_expiry(
243 &self,
244 id: &ObjectId,
245 target: ExpiryUpdate,
246 access_time: Timestamp,
247 ) -> Result<SetExpiryResponse>;
248
249 /// Deletes the object at the given path.
250 async fn delete_object(&self, id: &ObjectId, access_time: Timestamp) -> Result<DeleteResponse>;
251
252 /// Waits for any outstanding background operations to complete before shutdown.
253 ///
254 /// The default implementation is a no-op. Backends that spawn background tasks
255 /// (such as [`TieredStorage`](super::tiered::TieredStorage)) should override this
256 /// to wait for those tasks to complete.
257 async fn join(&self) {}
258
259 /// Borrows this backend as a [`MultipartUploadBackend`] if supported.
260 ///
261 /// The default returns an [`ErrorKind::Unsupported`]. Backends that implement
262 /// [`MultipartUploadBackend`] should override this to return `Ok(self)`.
263 fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> {
264 Err(ErrorKind::Unsupported.into())
265 }
266
267 /// Creates a resumable upload session for the object at `id`.
268 ///
269 /// Object metadata and its total length are declared upfront and cannot be mutated
270 /// during the upload.
271 ///
272 /// The returned string is opaque backend-defined state. [`StorageService`](crate::StorageService)
273 /// protects it before exposing the session token outside the service layer.
274 ///
275 /// Returns `Ok(None)` when this backend cannot store the described object resumably. Declining
276 /// is a routine outcome rather than an error, and the default implementation declines.
277 ///
278 /// # Errors
279 ///
280 /// Returns an error only when the backend supports resumable uploads but failed to open the
281 /// session.
282 async fn create_upload_session(
283 &self,
284 id: &ObjectId,
285 metadata: &Metadata,
286 upload_length: NonZeroU64,
287 ) -> Result<Option<BackendToken>> {
288 let _ = (id, metadata, upload_length);
289 Ok(None)
290 }
291
292 /// Writes a chunk of `content_length` bytes at `offset` into an open session.
293 ///
294 /// A backend may acknowledge fewer bytes than the chunk supplied, for example by persisting
295 /// only an aligned prefix. Callers must continue from the authoritative offset in the returned
296 /// [`UploadProgress`], or query [`Self::upload_offset`] after an ambiguous failure. A backend
297 /// may or may not accept a replay starting before its persisted offset.
298 ///
299 /// [`UploadProgress::Complete`] means the upload is terminal and the object is available
300 /// through this backend's normal read methods. A backend that composes another backend must
301 /// finish its own publication work before returning that outcome.
302 ///
303 /// A `content_length` of zero is valid. It writes nothing and reports the offset the backend
304 /// holds.
305 ///
306 /// Returns [`ErrorKind::UnknownUploadSession`] when `session` does not identify an open session,
307 /// and [`ErrorKind::ChunkExceedsUploadLength`] when the chunk would exceed the total length
308 /// declared when the session was created. Returns [`ErrorKind::ChunkTooSmall`] when a non-empty,
309 /// non-final chunk is shorter than the upload granularity.
310 async fn put_chunk(
311 &self,
312 session: &Session,
313 offset: u64,
314 content_length: u64,
315 stream: ClientStream,
316 ) -> Result<UploadProgress> {
317 let _ = (session, offset, content_length, stream);
318 Err(ErrorKind::Unsupported.into())
319 }
320
321 /// Reports how far the session has progressed.
322 ///
323 /// This can return [`UploadProgress::Complete`] repeatedly after the final chunk, including
324 /// when its original response was lost. A composed backend may finish pending idempotent
325 /// publication work before returning that terminal outcome.
326 ///
327 /// Returns [`ErrorKind::UnknownUploadSession`] when `session` does not identify a known session.
328 async fn upload_offset(&self, session: &Session) -> Result<UploadProgress> {
329 let _ = session;
330 Err(ErrorKind::Unsupported.into())
331 }
332
333 /// Cancels an upload session, discarding whatever was uploaded.
334 ///
335 /// Returns [`ErrorKind::UnknownUploadSession`] when `session` does not identify an open session.
336 async fn cancel_upload(&self, session: &Session) -> Result<()> {
337 let _ = session;
338 Err(ErrorKind::Unsupported.into())
339 }
340}
341
342/// Trait for backends that support our S3-style multipart upload protocol.
343#[async_trait::async_trait]
344pub trait MultipartUploadBackend: Backend + fmt::Debug + Send + Sync + 'static {
345 /// Initiates a new multipart upload at `id` with the given metadata.
346 async fn initiate_multipart(
347 &self,
348 id: &ObjectId,
349 metadata: &Metadata,
350 ) -> Result<InitiateMultipartResponse>;
351
352 /// Uploads a single part of the upload identified by `(id, upload_id)`.
353 async fn upload_part(
354 &self,
355 id: &ObjectId,
356 upload_id: &UploadId,
357 part_number: PartNumber,
358 content_length: u64,
359 content_md5: Option<&str>,
360 body: ClientStream,
361 ) -> Result<UploadPartResponse>;
362
363 /// Lists the parts uploaded so far for `(id, upload_id)`.
364 async fn list_parts(
365 &self,
366 id: &ObjectId,
367 upload_id: &UploadId,
368 max_parts: Option<u32>,
369 part_number_marker: Option<PartNumber>,
370 ) -> Result<ListPartsResponse>;
371
372 /// Aborts the upload identified by `(id, upload_id)`.
373 async fn abort_multipart(
374 &self,
375 id: &ObjectId,
376 upload_id: &UploadId,
377 ) -> Result<AbortMultipartResponse>;
378
379 /// Finalizes the upload identified by `(id, upload_id)` with the given
380 /// ordered list of parts.
381 ///
382 /// Note that this returns `Result<Option<CompleteMultipartError>>`.
383 /// It's therefore possible to get `Ok(Some(err))`, meaning that at the server level this will
384 /// translate to HTTP `200 OK` with an error contained in the response body.
385 /// We need to do it this way to mirror backends that also behave like this (namely S3 and
386 /// GCS).
387 async fn complete_multipart(
388 &self,
389 id: &ObjectId,
390 upload_id: &UploadId,
391 parts: Vec<CompletedPart>,
392 access_time: Timestamp,
393 ) -> Result<CompleteMultipartResponse>;
394}
395
396/// Trait for backends that support tombstone-conditional operations.
397///
398/// Only backends suitable for the high-volume tier of
399/// [`TieredStorage`](super::tiered::TieredStorage) implement this trait.
400/// The conditional methods provide atomic operations to avoid overwriting
401/// redirect tombstones.
402#[async_trait::async_trait]
403pub trait HighVolumeBackend: Backend {
404 /// Creates an upload marker separate from the logical object row.
405 ///
406 /// `revision` is the upload's unique LT revision. The fixed deadline is independent
407 /// of object expiration and must not be refreshed by subsequent requests.
408 async fn create_upload_marker(
409 &self,
410 revision: &ObjectId,
411 time_expires: Timestamp,
412 ) -> Result<()>;
413
414 /// Returns whether an upload marker exists and is live at `access_time`.
415 async fn has_upload_marker(&self, revision: &ObjectId, access_time: Timestamp) -> Result<bool>;
416
417 /// Atomically deletes a live upload marker.
418 ///
419 /// Returns `true` only when this call deletes a live marker.
420 /// Missing and expired markers return `false`, including on a repeated deletion.
421 async fn delete_upload_marker(
422 &self,
423 revision: &ObjectId,
424 access_time: Timestamp,
425 ) -> Result<bool>;
426
427 /// Writes the object only if NO redirect tombstone exists at this key.
428 ///
429 /// Returns `None` after storing the object, or `Some(tombstone)` (skipping
430 /// the write) when a redirect tombstone is present. The returned tombstone
431 /// carries the target LT `ObjectId` so the caller can route without a
432 /// second round trip.
433 ///
434 /// Takes [`Bytes`] instead of a [`ClientStream`] because callers on this
435 /// path have already fully buffered the payload.
436 async fn put_non_tombstone(
437 &self,
438 id: &ObjectId,
439 metadata: &Metadata,
440 payload: Bytes,
441 access_time: Timestamp,
442 ) -> Result<Option<Tombstone>>;
443
444 /// Retrieves (part of) an object with explicit tombstone awareness.
445 ///
446 /// Returns [`TieredGet::Tombstone`] instead of synthesizing a tombstone
447 /// object, making the caller's routing logic a compile-time distinction.
448 async fn get_tiered_object(
449 &self,
450 id: &ObjectId,
451 access_time: Timestamp,
452 range: Option<ByteRange>,
453 ) -> Result<TieredGet>;
454
455 /// Retrieves only metadata with explicit tombstone awareness.
456 ///
457 /// Implementations should skip the payload column where possible to avoid
458 /// fetching up to 1 MiB of data just to discover a tombstone.
459 async fn get_tiered_metadata(
460 &self,
461 id: &ObjectId,
462 access_time: Timestamp,
463 ) -> Result<TieredMetadata>;
464
465 /// Deletes the object only if it is NOT a redirect tombstone.
466 ///
467 /// Returns `None` after deleting the row (or if the row was already absent),
468 /// or `Some(tombstone)` (leaving the row intact) when the object is a
469 /// redirect tombstone. The returned tombstone carries the target LT
470 /// `ObjectId` so the caller can delete from long-term storage directly,
471 /// without a second round trip.
472 async fn delete_non_tombstone(
473 &self,
474 id: &ObjectId,
475 access_time: Timestamp,
476 ) -> Result<Option<Tombstone>>;
477
478 /// Atomically mutates the row if the current redirect state matches.
479 ///
480 /// `current` determines the precondition:
481 /// - `None`: succeeds only if no live tombstone exists (row absent, inline,
482 /// or tombstone present but logically expired).
483 /// - `Some(target)`: succeeds only if a tombstone exists whose redirect
484 /// resolves to `target`.
485 ///
486 /// **This operation is idempotent:** if the object is already in the target
487 /// state, it returns `true`. Whether the mutation runs again is up to the
488 /// implementation.
489 ///
490 /// Returns `true` on success or idempotent match, `false` if a conflicting
491 /// state was found (another writer won the race).
492 async fn compare_and_write(
493 &self,
494 id: &ObjectId,
495 current: Option<&ObjectId>,
496 write: TieredWrite,
497 access_time: Timestamp,
498 ) -> Result<bool>;
499
500 /// Atomically updates an existing row if its kind and redirect target match.
501 ///
502 /// `current = None` requires a live inline object. `Some(target)` requires
503 /// a live redirect to exactly that target. Updates never authorize creation
504 /// of an absent row.
505 ///
506 /// Returns [`SetExpiryResponse::Satisfied`] when applied or already satisfied,
507 /// [`SetExpiryResponse::NotFound`] when observed absent or expired, or
508 /// [`SetExpiryResponse::Rejected`] for an ineligible entry or failed condition.
509 /// Redirects can only resolve absolute targets because tombstones do not store
510 /// creation time. Backend failures are returned as errors.
511 async fn compare_and_update(
512 &self,
513 id: &ObjectId,
514 current: Option<&ObjectId>,
515 update: TieredUpdate,
516 access_time: Timestamp,
517 ) -> Result<SetExpiryResponse>;
518}
519
520/// Information about a redirect tombstone in the high-volume backend.
521#[derive(Clone, Debug, PartialEq, Eq)]
522pub struct Tombstone {
523 /// The [`ObjectId`] of the object in the long-term backend.
524 ///
525 /// For legacy tombstones with an empty `r` column, the HV backend resolves
526 /// this to the HV `ObjectId` itself before surfacing the tombstone to callers.
527 pub target: ObjectId,
528
529 /// The concrete deadline stored on the redirect.
530 pub time_expires: Option<Timestamp>,
531}
532
533impl Tombstone {
534 /// Returns whether the tombstone has expired at the given time.
535 pub fn is_expired(&self, now: Timestamp) -> bool {
536 self.time_expires.is_some_and(|deadline| deadline < now)
537 }
538}
539
540/// Typed response from [`HighVolumeBackend::get_tiered_object`].
541pub enum TieredGet {
542 /// A real object was found.
543 Object(Metadata, Option<ContentRange>, PayloadStream),
544 /// A redirect tombstone was found; the real object lives in the long-term backend.
545 Tombstone(Tombstone),
546 /// No entry exists at this key.
547 NotFound,
548}
549
550impl fmt::Debug for TieredGet {
551 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
552 match self {
553 TieredGet::Object(metadata, content_range, _stream) => f
554 .debug_tuple("Object")
555 .field(metadata)
556 .field(content_range)
557 .finish_non_exhaustive(),
558 TieredGet::Tombstone(info) => f.debug_tuple("Tombstone").field(info).finish(),
559 TieredGet::NotFound => write!(f, "NotFound"),
560 }
561 }
562}
563
564/// Typed metadata-only response from [`HighVolumeBackend::get_tiered_metadata`].
565#[derive(Debug)]
566pub enum TieredMetadata {
567 /// Metadata for a real object was found.
568 Object(Metadata),
569 /// A redirect tombstone was found; the real object lives in the long-term backend.
570 Tombstone(Tombstone),
571 /// No entry exists at this key.
572 NotFound,
573}
574
575/// The write operation performed by [`HighVolumeBackend::compare_and_write`].
576#[derive(Clone, Debug)]
577pub enum TieredWrite {
578 /// Write a redirect tombstone.
579 Tombstone(Tombstone),
580 /// Write inline object data.
581 Object(Metadata, Bytes),
582 /// Delete the row entirely.
583 Delete,
584}
585
586impl TieredWrite {
587 /// Returns the tombstone target if this is a tombstone write, or `None` otherwise.
588 pub fn target(&self) -> Option<&ObjectId> {
589 match self {
590 TieredWrite::Tombstone(t) => Some(&t.target),
591 _ => None,
592 }
593 }
594}
595
596/// The in-place operation performed by [`HighVolumeBackend::compare_and_update`].
597#[derive(Clone, Debug)]
598pub enum TieredUpdate {
599 /// Extend the deadline and adjust TTL duration while preserving other stored data.
600 SetExpiry(ExpiryUpdate),
601}
602
603/// Creates a reqwest client with required defaults.
604///
605/// Automatic decompression is disabled because backends store pre-compressed
606/// payloads and manage `Content-Encoding` themselves.
607pub(super) fn reqwest_client() -> reqwest::Client {
608 reqwest::Client::builder()
609 .user_agent(USER_AGENT)
610 .hickory_dns(true)
611 .http1_only()
612 .no_zstd()
613 .no_brotli()
614 .no_gzip()
615 .no_deflate()
616 .build()
617 // INVARIANT: Building fails only if the TLS backend cannot be initialized, which
618 // is checked at startup when the rustls crypto provider is installed.
619 .expect("failed to build backend HTTP client")
620}
621
622#[cfg(test)]
623mod tests {
624 use super::*;
625
626 #[test]
627 fn extended_ttl_rounding() {
628 let created = Timestamp::from_unix_secs(1_700_000_000).unwrap();
629 let policy = ExpirationPolicy::TimeToLive(Duration::from_hours(1));
630 for (seconds, rounded) in [
631 (0, 0),
632 (59, 59),
633 (61, 61),
634 (119, 120),
635 (121, 120),
636 (3598, 3600),
637 (3602, 3600),
638 (86_398, 86_400),
639 (86_402, 86_400),
640 (90 * 86_400 - 2, 90 * 86_400),
641 (90 * 86_400 + 2, 90 * 86_400),
642 (40_000, 39_600), // Exactly 1% rounds to hours.
643 (40_001, 40_020), // Just outside 1% falls back to minutes.
644 ] {
645 assert_eq!(
646 extended_expiration_policy(
647 policy,
648 Some(created),
649 created,
650 created + Duration::from_secs(seconds)
651 )
652 .unwrap(),
653 ExpirationPolicy::TimeToLive(Duration::from_secs(rounded)),
654 "{seconds} seconds",
655 );
656 }
657
658 assert_eq!(
659 round_duration(Duration::MAX),
660 Duration::from_secs(u64::MAX - u64::MAX % 86_400),
661 );
662 }
663
664 #[test]
665 fn extended_ttl_without_creation_time() {
666 // The inferred creation time may precede the epoch.
667 let epoch = Timestamp::from_unix_secs(0).unwrap();
668 let policy = ExpirationPolicy::TimeToLive(Duration::from_hours(1));
669 let first_expiry = epoch + Duration::from_hours(1);
670 let extended = extended_expiration_policy(policy, None, epoch, first_expiry).unwrap();
671 assert_eq!(
672 extended,
673 ExpirationPolicy::TimeToLive(Duration::from_hours(2))
674 );
675
676 let next_expiry = first_expiry + Duration::from_hours(1);
677 let extended_again =
678 extended_expiration_policy(extended, None, first_expiry, next_expiry).unwrap();
679 assert_eq!(
680 extended_again,
681 ExpirationPolicy::TimeToLive(Duration::from_hours(3))
682 );
683 }
684
685 #[test]
686 fn extended_ttl_fractional_duration() {
687 let epoch = Timestamp::from_unix_secs(0).unwrap();
688 assert_eq!(
689 extended_expiration_policy(
690 ExpirationPolicy::TimeToLive(Duration::from_millis(1500)),
691 None,
692 epoch,
693 epoch + Duration::from_secs(1)
694 )
695 .unwrap(),
696 ExpirationPolicy::TimeToLive(Duration::from_secs(3)),
697 );
698 }
699
700 #[test]
701 fn extended_ttl_invalid_creation_time() {
702 let expiry = Timestamp::from_unix_secs(1_700_000_000).unwrap();
703 let created = expiry + Duration::from_secs(1);
704 let policy = ExpirationPolicy::TimeToLive(Duration::from_hours(1));
705 assert_eq!(
706 extended_expiration_policy(policy, Some(created), expiry, expiry)
707 .unwrap_err()
708 .kind(),
709 ErrorKind::CorruptData
710 );
711 }
712
713 #[test]
714 fn extended_ttl_duration_overflow() {
715 let epoch = Timestamp::from_unix_secs(0).unwrap();
716 assert_eq!(
717 extended_expiration_policy(
718 ExpirationPolicy::TimeToLive(Duration::MAX),
719 None,
720 epoch,
721 epoch + Duration::from_secs(1),
722 )
723 .unwrap_err()
724 .kind(),
725 ErrorKind::CorruptData
726 );
727 }
728
729 #[test]
730 fn extended_ttl_saturated_deadline() {
731 let max = Timestamp::from_unix_secs(253_402_300_799).unwrap();
732 assert_eq!(
733 extended_expiration_policy(
734 ExpirationPolicy::TimeToLive(Duration::from_hours(3)),
735 Some(max - Duration::from_hours(1)),
736 max - Duration::from_secs(1),
737 max.saturating_add(Duration::from_hours(1))
738 )
739 .unwrap(),
740 ExpirationPolicy::TimeToLive(Duration::from_hours(1))
741 );
742 }
743
744 #[test]
745 fn expiry_target_resolution() {
746 let created = Timestamp::from_unix_secs(1_700_000_000).unwrap();
747 assert_eq!(
748 ExpiryTarget::FromCreation(Duration::ZERO).resolve(Some(created)),
749 Some(created)
750 );
751 assert_eq!(ExpiryTarget::At(created).resolve(None), Some(created));
752 assert_eq!(
753 ExpiryTarget::FromCreation(Duration::ZERO).resolve(None),
754 None
755 );
756 let max = Timestamp::from_unix_secs(253_402_300_799).unwrap();
757 assert_eq!(
758 ExpiryTarget::FromCreation(Duration::from_secs(1)).resolve(Some(max)),
759 Some(max)
760 );
761
762 // The cap is measured from access time, not creation.
763 let deadline = created + Duration::from_hours(3);
764 let access_time = created + Duration::from_hours(1);
765 for target in [
766 ExpiryTarget::At(deadline),
767 ExpiryTarget::FromCreation(Duration::from_hours(3)),
768 ] {
769 let update = ExpiryUpdate {
770 target,
771 max: Some(Duration::from_hours(1)),
772 };
773 assert_eq!(
774 update
775 .resolve(Some(created), access_time)
776 .unwrap_err()
777 .kind(),
778 ErrorKind::InvalidMetadata
779 );
780 assert_eq!(
781 update
782 .resolve(Some(created), access_time + Duration::from_hours(1))
783 .unwrap(),
784 Some(deadline)
785 );
786 }
787 }
788}