Skip to main content

loonfs_objectstore/
immutable_write.rs

1//! Verified writes for objects whose keys name immutable bytes.
2
3use crate::retry::{
4    transport_retry_backoff, transport_retry_pause, OperationDeadline, TransportRetryPolicy,
5};
6use crate::timing::StdMonotonicTimer;
7use crate::{
8    ObjectMetadata, ObjectStore, ObjectStoreError, PutMode, PROVIDER_MULTIPART_THRESHOLD_BYTES,
9};
10use bytes::Bytes;
11use thiserror::Error;
12
13/// Failure to establish that an immutable key contains the requested bytes.
14#[derive(Debug, Error)]
15#[non_exhaustive]
16pub enum ImmutableWriteError {
17    /// The key names bytes other than the immutable payload supplied here.
18    #[error("immutable object `{object_key}` already exists with different bytes")]
19    DifferentObject {
20        /// Durable key whose existing bytes violated immutability.
21        object_key: String,
22    },
23    /// The storage boundary failed before byte identity could be established.
24    #[error("immutable write transport failed for `{object_key}`: {source}")]
25    Transport {
26        /// Durable key whose byte identity could not be established.
27        object_key: String,
28        /// Final storage failure after retry and read-back reconciliation.
29        #[source]
30        source: ObjectStoreError,
31    },
32}
33
34impl ImmutableWriteError {
35    /// Returns the immutable object key whose byte identity was not established.
36    pub fn object_key(&self) -> &str {
37        match self {
38            Self::DifferentObject { object_key } | Self::Transport { object_key, .. } => object_key,
39        }
40    }
41}
42
43pub(crate) async fn put<S: ObjectStore + ?Sized>(
44    store: &S,
45    key: &str,
46    bytes: Bytes,
47) -> std::result::Result<(), ImmutableWriteError> {
48    let timer = StdMonotonicTimer::default();
49    let retry_policy = TransportRetryPolicy::DEFAULT;
50    let mode = if bytes.len() as u64 >= PROVIDER_MULTIPART_THRESHOLD_BYTES {
51        PutMode::Overwrite
52    } else {
53        PutMode::CreateIfAbsent
54    };
55    let deadline = OperationDeadline::start(&timer, retry_policy.operation_deadline);
56    let mut retries = 0;
57    let mut ambiguous_transport = None;
58
59    loop {
60        match store.put(key, bytes.clone(), mode.clone()).await {
61            Ok(_) => return Ok(()),
62            Err(error @ ObjectStoreError::PreconditionFailed { .. }) => {
63                return resolve_readback(store, key, &bytes, ambiguous_transport.unwrap_or(error))
64                    .await;
65            }
66            Err(error @ ObjectStoreError::Transport { .. }) => {
67                let Some(remaining) = deadline.remaining() else {
68                    return resolve_readback(store, key, &bytes, error).await;
69                };
70                if retries >= retry_policy.max_retries {
71                    return resolve_readback(store, key, &bytes, error).await;
72                }
73                retries += 1;
74                let backoff = transport_retry_backoff(&retry_policy, retries).min(remaining);
75                tracing::info!(
76                    object_key = key,
77                    operation = "put_immutable_verified",
78                    retry = retries,
79                    max_retries = retry_policy.max_retries,
80                    backoff_ms = u64::try_from(backoff.as_millis()).unwrap_or(u64::MAX),
81                    error = %error,
82                    "transient immutable write failure, backing off before retry",
83                );
84                ambiguous_transport = Some(error);
85                transport_retry_pause(backoff).await;
86            }
87            Err(error) if ambiguous_transport.is_some() => {
88                return resolve_readback(store, key, &bytes, error).await;
89            }
90            Err(error) => return Err(transport(key, error)),
91        }
92    }
93}
94
95async fn resolve_readback<S: ObjectStore + ?Sized>(
96    store: &S,
97    key: &str,
98    expected: &Bytes,
99    original: ObjectStoreError,
100) -> std::result::Result<(), ImmutableWriteError> {
101    match readback(store, key, expected).await {
102        Ok(ImmutableReadback::Identical(_)) => Ok(()),
103        Ok(ImmutableReadback::Different) => Err(ImmutableWriteError::DifferentObject {
104            object_key: key.to_owned(),
105        }),
106        Ok(ImmutableReadback::Missing) => Err(transport(key, original)),
107        Err(verify_error @ ObjectStoreError::Transport { .. }) => {
108            let source = ObjectStoreError::transport(
109                key,
110                format!(
111                    "{}; failed to verify immutable write outcome: {verify_error}",
112                    original.message()
113                ),
114            );
115            Err(transport(key, source))
116        }
117        Err(verify_error) => Err(transport(key, verify_error)),
118    }
119}
120
121pub(crate) enum ImmutableReadback {
122    Identical(ObjectMetadata),
123    Different,
124    Missing,
125}
126
127pub(crate) async fn readback<S: ObjectStore + ?Sized>(
128    store: &S,
129    key: &str,
130    expected: &Bytes,
131) -> crate::Result<ImmutableReadback> {
132    match store.get_with_metadata(key).await? {
133        Some(body) if body.bytes.as_slice() == expected.as_ref() => {
134            Ok(ImmutableReadback::Identical(body.metadata))
135        }
136        Some(_) => Ok(ImmutableReadback::Different),
137        None => Ok(ImmutableReadback::Missing),
138    }
139}
140
141fn transport(key: &str, source: ObjectStoreError) -> ImmutableWriteError {
142    ImmutableWriteError::Transport {
143        object_key: key.to_owned(),
144        source,
145    }
146}
147
148#[cfg(test)]
149mod tests {
150    use super::*;
151    use crate::test_support::SteppingTimer;
152    use std::time::Duration;
153
154    #[test]
155    fn immutable_retry_deadline_is_one_budget_across_attempts() {
156        let timer = SteppingTimer::new(70_000);
157        let deadline = OperationDeadline::start(&timer, Duration::from_secs(120));
158        assert_eq!(deadline.remaining(), Some(Duration::from_secs(50)));
159        assert_eq!(deadline.remaining(), None);
160    }
161}