qubit_fs/write/
async_write_all_operation.rs1use std::fmt::Debug;
10use std::fmt::Formatter;
11use std::fmt::Result as FmtResult;
12use std::io::Error as IoError;
13
14use qubit_io::AsyncOutput;
15
16use crate::AsyncFileSystem;
17use crate::error::FsError;
18use crate::error::FsErrorKind;
19use crate::error::FsOperation;
20use crate::error::OpenFailureStage;
21use crate::metadata::WriteOutcome;
22use crate::path::Path;
23use crate::write::AsyncWriteAllOperationFailure;
24use crate::write::AsyncWriteAllOperationState;
25use crate::write::AsyncWriterRecovery;
26use crate::write::WriteFailureState;
27use crate::write::WriteOptions;
28use crate::write::WriterState;
29use crate::write::internal::WriteAllCancellationGuard;
30use crate::write::internal::WriteAllRecoverySnapshot;
31use crate::write::internal::open_failure_state;
32#[must_use]
59pub struct AsyncWriteAllOperation {
60 filesystem: AsyncFileSystem,
62 path: Path,
64 bytes: Vec<u8>,
66 options: WriteOptions,
68 state: AsyncWriteAllOperationState,
70 writer: Option<AsyncWriterRecovery>,
72 recovery: WriteAllRecoverySnapshot,
74}
75impl AsyncWriteAllOperation {
76 pub(crate) fn new(filesystem: AsyncFileSystem, path: Path, bytes: Vec<u8>, options: WriteOptions) -> Self {
78 Self {
79 filesystem,
80 path,
81 bytes,
82 options,
83 state: AsyncWriteAllOperationState::Ready,
84 writer: None,
85 recovery: WriteAllRecoverySnapshot::new(),
86 }
87 }
88 #[inline]
90 #[must_use]
91 pub const fn filesystem(&self) -> &AsyncFileSystem {
92 &self.filesystem
93 }
94 #[inline]
96 #[must_use]
97 pub const fn path(&self) -> &Path {
98 &self.path
99 }
100 #[inline]
102 #[must_use = "inspect publication state before choosing a recovery action"]
103 pub const fn state(&self) -> AsyncWriteAllOperationState {
104 self.state
105 }
106 #[inline]
108 #[must_use]
109 pub const fn has_recovery(&self) -> bool {
110 self.writer.is_some()
111 }
112 #[inline]
114 #[must_use]
115 pub fn recovery(&mut self) -> Option<&mut AsyncWriterRecovery> {
116 self.writer.as_mut()
117 }
118 #[inline]
123 #[must_use]
124 pub fn take_recovery(&mut self) -> Option<AsyncWriterRecovery> {
125 self.writer.take()
126 }
127 #[inline]
129 #[must_use]
130 pub const fn written_bytes(&self) -> u64 {
131 self.recovery.written_bytes
132 }
133 pub async fn execute(&mut self) -> Result<WriteOutcome, AsyncWriteAllOperationFailure> {
154 if self.state != AsyncWriteAllOperationState::Ready {
155 return Err(AsyncWriteAllOperationFailure::new(
156 invalid_state(&self.path, &self.filesystem),
157 self.recovery.state,
158 self.recovery.written_bytes,
159 ));
160 }
161 let Self {
162 filesystem,
163 path,
164 bytes,
165 options,
166 state,
167 writer,
168 recovery,
169 } = self;
170 let bytes = std::mem::take(bytes);
171 let mut guard = WriteAllCancellationGuard::start(state, writer, recovery);
172 let result = execute_write(filesystem, path, &bytes, options, guard.writer_mut()).await;
173 guard.finish(&result);
174 result
175 }
176}
177async fn execute_write(
179 filesystem: &AsyncFileSystem,
180 path: &Path,
181 bytes: &[u8],
182 options: &WriteOptions,
183 slot: &mut Option<AsyncWriterRecovery>,
184) -> Result<WriteOutcome, AsyncWriteAllOperationFailure> {
185 if slot.is_none() {
186 match filesystem.open_writer(path, options.clone()).await {
187 Ok(writer) => *slot = Some(AsyncWriterRecovery::Opened(Box::new(writer))),
188 Err(failure) => {
189 let (error, stage, recovery) = failure.into_parts();
190 *slot = recovery.map(AsyncWriterRecovery::Rejected);
191 let state = match stage {
192 OpenFailureStage::Preflight => WriteFailureState::NotPublished,
193 OpenFailureStage::ProviderOpen => open_failure_state(&error),
194 OpenFailureStage::OutcomeValidation => WriteFailureState::Indeterminate,
195 };
196 return Err(AsyncWriteAllOperationFailure::new(error, state, 0));
197 }
198 }
199 }
200 let writer = slot
201 .as_mut()
202 .and_then(AsyncWriterRecovery::opened_mut)
203 .expect("writer is retained after open");
204 if let Err(error) = writer.write_fully_async(bytes).await {
205 let error = contextual(filesystem, error, path);
206 let state = state_for(error.has_indeterminate_effect(), writer.state());
207 return Err(AsyncWriteAllOperationFailure::new(error, state, writer.written_bytes()));
208 }
209 if let Err(error) = writer.flush_async().await {
210 let error = contextual(filesystem, error, path);
211 let state = state_for(error.has_indeterminate_effect(), writer.state());
212 return Err(AsyncWriteAllOperationFailure::new(error, state, writer.written_bytes()));
213 }
214 match writer.commit_async().await {
215 Ok(outcome) => Ok(outcome),
216 Err(failure) => {
217 let (error, state) = failure.into_parts();
218 Err(AsyncWriteAllOperationFailure::new(error, state, writer.written_bytes()))
219 }
220 }
221}
222fn contextual(filesystem: &AsyncFileSystem, error: IoError, path: &Path) -> FsError {
224 filesystem.core().enrich(
225 FsError::from_stream_io(error, FsOperation::Write, path),
226 Some(path),
227 FsOperation::Write,
228 )
229}
230fn state_for(indeterminate: bool, state: WriterState) -> WriteFailureState {
232 if indeterminate {
233 WriteFailureState::Indeterminate
234 } else {
235 state.publication_failure_state()
236 }
237}
238fn invalid_state(path: &Path, filesystem: &AsyncFileSystem) -> FsError {
240 FsError::new(
241 FsErrorKind::InvalidState,
242 FsOperation::Write,
243 "async whole-file write cannot execute in its current state",
244 )
245 .with_path(path.clone())
246 .with_provider(filesystem.properties().info().provider_id())
247}
248
249impl Debug for AsyncWriteAllOperation {
250 fn fmt(&self, formatter: &mut Formatter<'_>) -> FmtResult {
252 formatter
253 .debug_struct("AsyncWriteAllOperation")
254 .field("path", &self.path)
255 .field("state", &self.state)
256 .field("written_bytes", &self.recovery.written_bytes)
257 .field("has_recovery", &self.writer.is_some())
258 .finish()
259 }
260}