loonfs_objectstore/
immutable_write.rs1use 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#[derive(Debug, Error)]
15#[non_exhaustive]
16pub enum ImmutableWriteError {
17 #[error("immutable object `{object_key}` already exists with different bytes")]
19 DifferentObject {
20 object_key: String,
22 },
23 #[error("immutable write transport failed for `{object_key}`: {source}")]
25 Transport {
26 object_key: String,
28 #[source]
30 source: ObjectStoreError,
31 },
32}
33
34impl ImmutableWriteError {
35 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}