1use std::pin::Pin;
12
13use crate::AsyncFileSystem;
14use crate::error::FsError;
15use crate::error::FsErrorKind;
16use crate::error::FsOperation;
17use crate::error::FsResult;
18use crate::metadata::AchievedAtomicity;
19use crate::metadata::AtomicityRequirement;
20use crate::path::Path;
21use crate::path::PathComponent;
22use crate::spi::AsyncTempResourceSpi;
23use crate::spi::PersistRequest;
24use crate::spi::SpiFuture;
25use crate::temp::PersistFailure;
26use crate::temp::PersistFailureState;
27use crate::temp::PersistOptions;
28use crate::temp::PersistOutcome;
29use crate::temp::TempResourceState;
30use crate::temp::internal::TempLifecycle;
31
32pub struct AsyncTempFile {
53 file_system: AsyncFileSystem,
55 path: Path,
57 session: Pin<Box<dyn AsyncTempResourceSpi>>,
59 lifecycle: TempLifecycle,
61 resource_name: &'static str,
63}
64
65impl AsyncTempFile {
66 pub(crate) fn new(
76 file_system: AsyncFileSystem,
77 path: Path,
78 session: Box<dyn AsyncTempResourceSpi>,
79 resource_name: &'static str,
80 ) -> Self {
81 Self {
82 file_system,
83 path,
84 session: Box::into_pin(session),
85 lifecycle: TempLifecycle::new(),
86 resource_name,
87 }
88 }
89
90 #[inline]
95 #[must_use]
96 pub const fn path(&self) -> &Path {
97 &self.path
98 }
99
100 #[inline]
105 #[must_use]
106 pub const fn state(&self) -> TempResourceState {
107 self.lifecycle.state()
108 }
109
110 #[inline]
112 #[must_use]
113 pub fn child(&self, component: &PathComponent) -> Path {
114 self.path.child(component)
115 }
116
117 #[inline]
119 #[must_use]
120 pub fn descendant(&self, relative: &crate::path::RelativePath) -> Path {
121 self.path.join(relative)
122 }
123
124 #[inline]
133 pub fn cleanup(&mut self) -> SpiFuture<'_, FsResult<()>> {
134 self.lifecycle("cannot be cleaned now", FsOperation::CleanupTemp, |session| {
135 session.cleanup()
136 })
137 }
138
139 #[inline]
149 pub fn keep(&mut self) -> SpiFuture<'_, Result<PersistOutcome, PersistFailure>> {
150 if self.lifecycle.state() != TempResourceState::Owned {
151 let error = self.invalid_state(FsOperation::KeepTemp, "cannot be kept now");
152 return Box::pin(async move {
153 Err(PersistFailure::new(error, self.lifecycle.failure_state())
154 .with_publication_target(self.lifecycle.publication_target()))
155 });
156 }
157 Box::pin(async move {
158 self.lifecycle.begin_pending();
159 match self.session.as_mut().keep().await {
160 Ok(outcome) => {
161 if let Err(error) = self.file_system.validate_temp_keep_target(&self.path, outcome.target()) {
162 return Err(PersistFailure::new(error, PersistFailureState::Indeterminate));
163 }
164 self.path = outcome.target().clone();
165 self.lifecycle.record_success(true, outcome.target().clone());
166 Ok(outcome)
167 }
168 Err(failure) => {
169 let (error, state) = failure.into_parts();
170 self.lifecycle.record_failure(state, error.target().cloned(), true);
171 let target = self.path.clone();
172 Err(PersistFailure::new(
173 error.with_operation(FsOperation::KeepTemp).with_missing_context(
174 &self.path,
175 Some(&target),
176 self.file_system.properties().info().provider_id(),
177 ),
178 state,
179 )
180 .with_publication_target(self.lifecycle.publication_target()))
181 }
182 }
183 })
184 }
185
186 pub fn persist<'a>(
201 &'a mut self,
202 target: &'a Path,
203 options: PersistOptions,
204 ) -> SpiFuture<'a, Result<PersistOutcome, PersistFailure>> {
205 if self.lifecycle.state() != TempResourceState::Owned {
206 let error = self.invalid_state(FsOperation::PersistTemp, "cannot be persisted now");
207 return Box::pin(async move {
208 Err(PersistFailure::new(error, self.lifecycle.failure_state())
209 .with_publication_target(self.lifecycle.publication_target()))
210 });
211 }
212 if let Err(error) = self.file_system.preflight_temp_persist(&self.path, target, &options) {
213 return Box::pin(async move { Err(PersistFailure::new(error, PersistFailureState::NotPublished)) });
214 }
215 Box::pin(async move {
216 self.lifecycle.begin_pending();
217 let atomicity = options.atomicity();
218 let result = self
219 .session
220 .as_mut()
221 .persist(PersistRequest::new(target, options))
222 .await;
223 match &result {
224 Ok(outcome)
225 if outcome.target() == target
226 && !(atomicity == AtomicityRequirement::Required
227 && outcome.atomicity() != AchievedAtomicity::Atomic) =>
228 {
229 self.lifecycle.record_success(false, outcome.target().clone());
230 }
231 Ok(outcome) if outcome.target() != target => {
232 self.lifecycle
233 .record_failure(PersistFailureState::Indeterminate, Some(target.clone()), false);
234 }
235 Ok(_) => self.lifecycle.record_failure(
236 PersistFailureState::PublishedSourceRetained,
237 Some(target.clone()),
238 false,
239 ),
240 Err(failure) => self
241 .lifecycle
242 .record_failure(failure.state(), Some(target.clone()), false),
243 }
244 match result {
245 Ok(outcome) if outcome.target() != target => Err(PersistFailure::new(
246 FsError::new(
247 FsErrorKind::ProviderContractViolation,
248 FsOperation::PersistTemp,
249 "provider reported a persistence target different from the request",
250 )
251 .with_path(self.path.clone())
252 .with_target(target.clone()),
253 PersistFailureState::Indeterminate,
254 )),
255 Ok(outcome)
256 if atomicity == AtomicityRequirement::Required
257 && outcome.atomicity() != AchievedAtomicity::Atomic =>
258 {
259 Err(PersistFailure::new(
260 FsError::new(
261 FsErrorKind::ProviderContractViolation,
262 FsOperation::PersistTemp,
263 "provider reported non-atomic success for atomic-required persist",
264 )
265 .with_path(self.path.clone())
266 .with_target(target.clone()),
267 PersistFailureState::PublishedSourceRetained,
268 )
269 .with_publication_target(self.lifecycle.publication_target()))
270 }
271 Err(failure) => {
272 let (error, state) = failure.into_parts();
273 Err(PersistFailure::new(self.contextual_persist_error(error, target), state)
274 .with_publication_target(self.lifecycle.publication_target()))
275 }
276 Ok(outcome) => Ok(outcome),
277 }
278 })
279 }
280
281 fn lifecycle<'a, F>(
299 &'a mut self,
300 action: &'static str,
301 operation: FsOperation,
302 call: F,
303 ) -> SpiFuture<'a, FsResult<()>>
304 where
305 F: FnOnce(Pin<&'a mut dyn AsyncTempResourceSpi>) -> SpiFuture<'a, FsResult<()>> + Send + 'a,
306 {
307 if !matches!(
308 self.lifecycle.state(),
309 TempResourceState::Owned | TempResourceState::CleanupRequired
310 ) {
311 let error = self.invalid_state(operation, action);
312 return Box::pin(async move { Err(error) });
313 }
314 Box::pin(async move {
315 let previous_lifecycle = self.lifecycle.clone();
316 self.lifecycle.begin_pending();
317 let result = call(self.session.as_mut()).await;
318 self.lifecycle = previous_lifecycle;
319 match &result {
320 Ok(()) => self.lifecycle.record_cleanup_success(),
321 Err(error) => self.lifecycle.record_cleanup_error(error),
322 }
323 result.map_err(|error| {
324 error.with_operation(operation).with_missing_context(
325 &self.path,
326 None,
327 self.file_system.properties().info().provider_id(),
328 )
329 })
330 })
331 }
332
333 fn invalid_state(&self, operation: FsOperation, action: &str) -> FsError {
342 let message = format!("{} {}", self.resource_name, action);
343 FsError::new(FsErrorKind::InvalidState, operation, &message).with_path(self.path.clone())
344 }
345
346 fn contextual_persist_error(&self, error: FsError, target: &Path) -> FsError {
356 error.with_operation(FsOperation::PersistTemp).with_missing_context(
357 &self.path,
358 Some(target),
359 self.file_system.properties().info().provider_id(),
360 )
361 }
362}
363
364impl Drop for AsyncTempFile {
365 fn drop(&mut self) {
366 if matches!(
367 self.lifecycle.state(),
368 TempResourceState::Owned | TempResourceState::CleanupRequired
369 ) {
370 self.session.as_mut().cancel_on_drop();
371 }
372 }
373}