1use std::fmt::Debug;
12use std::fmt::Formatter;
13use std::fmt::Result as FmtResult;
14use std::io::Error as IoError;
15use std::io::ErrorKind as IoErrorKind;
16use std::io::Result as IoResult;
17use std::pin::Pin;
18use std::task::Context;
19use std::task::Poll;
20
21use qubit_io::AsyncOutput;
22
23use crate::error::FsEffectState;
24use crate::error::FsError;
25use crate::error::FsErrorKind;
26use crate::error::FsOperation;
27use crate::facade::facade_core::FacadeCore;
28use crate::facade::internal::ByteBudget;
29use crate::facade::internal::FileSystemResource;
30use crate::metadata::AchievedAtomicity;
31use crate::metadata::AtomicityRequirement;
32use crate::metadata::DurabilityRequirement;
33use crate::metadata::OpenedFileInfo;
34use crate::metadata::WriteOutcome;
35use crate::spi::AsyncFileWriteSession;
36use crate::spi::SpiFuture;
37use crate::write::WriteAbortOutcome;
38use crate::write::WriteFailure;
39use crate::write::WriteFailureState;
40use crate::write::WriterState;
41
42pub struct AsyncFileWriter {
66 session: Pin<Box<dyn AsyncFileWriteSession>>,
68 info: OpenedFileInfo,
70 state: WriterState,
72 abort_completed: bool,
74 atomicity: AtomicityRequirement,
76 durability: DurabilityRequirement,
78 provider: Box<str>,
80 write_budget: Option<ByteBudget>,
82 written_bytes: u64,
84}
85
86impl AsyncFileWriter {
87 #[inline]
96 #[must_use]
97 pub(crate) fn new(
98 info: OpenedFileInfo,
99 session: Box<dyn AsyncFileWriteSession>,
100 atomicity: AtomicityRequirement,
101 durability: DurabilityRequirement,
102 provider: &str,
103 max_write_bytes: Option<u64>,
104 ) -> Self {
105 Self {
106 session: Box::into_pin(session),
107 info,
108 state: WriterState::Open,
109 abort_completed: false,
110 atomicity,
111 durability,
112 provider: provider.into(),
113 write_budget: max_write_bytes
114 .map(|maximum| FacadeCore::byte_budget(FileSystemResource::WriteBytes, maximum)),
115 written_bytes: 0,
116 }
117 }
118
119 #[inline]
124 #[must_use]
125 pub fn info(&self) -> &OpenedFileInfo {
126 &self.info
127 }
128
129 #[inline]
134 #[must_use]
135 pub const fn state(&self) -> WriterState {
136 self.state
137 }
138
139 #[inline]
141 #[must_use]
142 pub(crate) const fn written_bytes(&self) -> u64 {
143 self.written_bytes
144 }
145
146 #[inline]
148 pub(crate) fn mark_indeterminate(&mut self) {
149 self.state = WriterState::Indeterminate;
150 }
151
152 pub fn commit_async(&mut self) -> SpiFuture<'_, Result<WriteOutcome, WriteFailure>> {
164 if self.state != WriterState::Open {
165 let publication_state = self.state.publication_failure_state();
166 let error = self.invalid_state(
167 FsOperation::CommitWriter,
168 "writer cannot be committed in its current state",
169 );
170 return Box::pin(async move { Err(WriteFailure::new(error, publication_state)) });
171 }
172 Box::pin(async move {
173 self.state = WriterState::Indeterminate;
174 let result = self.session.as_mut().commit_async().await;
175 match result {
176 Ok(outcome) => {
177 self.state = WriterState::Committed;
178 if self.atomicity == AtomicityRequirement::Required
179 && outcome.atomicity() != AchievedAtomicity::Atomic
180 {
181 self.state = WriterState::Published;
182 return Err(WriteFailure::new(
183 FsError::new(
184 FsErrorKind::ProviderContractViolation,
185 FsOperation::CommitWriter,
186 "provider reported non-atomic success for an atomic-required write",
187 )
188 .with_path(self.info.path().clone())
189 .with_provider(&self.provider)
190 .with_effect_state(FsEffectState::Applied),
191 WriteFailureState::Published,
192 ));
193 }
194 if self.durability == DurabilityRequirement::Required && !outcome.durable() {
195 self.state = WriterState::Published;
196 return Err(WriteFailure::new(
197 FsError::new(
198 FsErrorKind::ProviderContractViolation,
199 FsOperation::CommitWriter,
200 "provider reported non-durable success for a durability-required write",
201 )
202 .with_path(self.info.path().clone())
203 .with_provider(&self.provider)
204 .with_effect_state(FsEffectState::Applied),
205 WriteFailureState::Published,
206 ));
207 }
208 if let Some(bytes_written) = outcome.bytes_written()
209 && bytes_written != self.written_bytes
210 {
211 self.state = WriterState::Published;
212 return Err(WriteFailure::new(
213 FsError::new(
214 FsErrorKind::ProviderContractViolation,
215 FsOperation::CommitWriter,
216 "provider reported a byte count different from the bytes accepted by the writer",
217 )
218 .with_path(self.info.path().clone())
219 .with_provider(&self.provider)
220 .with_effect_state(FsEffectState::Applied),
221 WriteFailureState::Published,
222 ));
223 }
224 Ok(outcome)
225 }
226 Err(failure) => {
227 self.state = match failure.state() {
228 WriteFailureState::RetryableNotPublished => WriterState::Open,
229 WriteFailureState::NotPublished => WriterState::NotPublished,
230 WriteFailureState::Published => WriterState::Published,
231 WriteFailureState::Indeterminate => WriterState::Indeterminate,
232 };
233 let (error, state) = failure.into_parts();
234 Err(WriteFailure::new(
235 self.contextual_error(error, FsOperation::CommitWriter),
236 state,
237 ))
238 }
239 }
240 })
241 }
242
243 pub fn abort_async(&mut self) -> SpiFuture<'_, crate::error::FsResult<WriteAbortOutcome>> {
254 if self.abort_completed
255 || !matches!(
256 self.state,
257 WriterState::Open | WriterState::NotPublished | WriterState::Published | WriterState::Indeterminate
258 )
259 {
260 let error = self.invalid_state(
261 FsOperation::AbortWriter,
262 "writer cannot be aborted in its current state",
263 );
264 return Box::pin(async move { Err(error) });
265 }
266 Box::pin(async move {
267 let previous_state = self.state;
268 self.state = WriterState::Indeterminate;
269 match self.session.as_mut().abort_async().await {
270 Ok(outcome) => {
271 self.abort_completed = true;
272 self.state = match outcome {
273 WriteAbortOutcome::NotPublished => WriterState::Aborted,
274 WriteAbortOutcome::Published => WriterState::Published,
275 WriteAbortOutcome::Indeterminate => WriterState::Indeterminate,
276 };
277 Ok(outcome)
278 }
279 Err(error) => {
280 if !error.has_indeterminate_effect() {
281 self.state = previous_state;
282 }
283 Err(self.contextual_error(error, FsOperation::AbortWriter))
284 }
285 }
286 })
287 }
288
289 fn invalid_state(&self, operation: FsOperation, message: &str) -> FsError {
291 FsError::new(FsErrorKind::InvalidState, operation, message)
292 .with_path(self.info.path().clone())
293 .with_provider(&self.provider)
294 }
295
296 fn closed_io_error(&self) -> IoError {
298 IoError::new(
299 IoErrorKind::BrokenPipe,
300 self.invalid_state(FsOperation::Write, "writer no longer accepts bytes"),
301 )
302 }
303
304 fn check_write_limit(&self, count: usize) -> IoResult<u64> {
306 let count = FacadeCore::quantity_from_usize(count, FsOperation::Write, self.info.path(), &self.provider)
307 .map_err(FsError::into_io_error)?;
308 if let Some(budget) = &self.write_budget
309 && let Err(error) = budget.check_available(count)
310 {
311 return Err(FacadeCore::budget_error(
312 error,
313 FsOperation::Write,
314 self.info.path(),
315 &self.provider,
316 "write session exceeds the provider byte limit",
317 )
318 .into_io_error());
319 }
320 Ok(count)
321 }
322
323 fn record_written_bytes(&mut self, count: usize) -> IoResult<()> {
329 let count = FacadeCore::quantity_from_usize(count, FsOperation::Write, self.info.path(), &self.provider)
330 .map_err(FsError::into_io_error)?;
331 if let Some(error) = self
332 .write_budget
333 .as_mut()
334 .and_then(|budget| budget.try_consume(count).err())
335 {
336 return Err(FacadeCore::budget_error(
337 error,
338 FsOperation::Write,
339 self.info.path(),
340 &self.provider,
341 "write session exceeds the provider byte limit",
342 )
343 .into_io_error());
344 }
345 self.written_bytes = self
346 .written_bytes
347 .checked_add(count)
348 .ok_or_else(|| self.byte_count_error())?;
349 Ok(())
350 }
351
352 fn byte_count_error(&self) -> IoError {
355 FsError::new(
356 FsErrorKind::ResourceLimitExceeded,
357 FsOperation::Write,
358 "write byte count exceeds the filesystem API reporting range",
359 )
360 .with_path(self.info.path().clone())
361 .with_provider(&self.provider)
362 .into_io_error()
363 }
364
365 fn contextual_error(&self, error: FsError, operation: FsOperation) -> FsError {
367 error
368 .with_operation(operation)
369 .with_missing_context(self.info.path(), None, &self.provider)
370 }
371}
372
373impl AsyncOutput for AsyncFileWriter {
374 type Item = u8;
375
376 #[inline]
377 fn is_buffered(&self) -> bool {
378 self.session.is_buffered()
379 }
380
381 unsafe fn poll_write_unchecked(
382 self: Pin<&mut Self>,
383 cx: &mut Context<'_>,
384 input: &[u8],
385 index: usize,
386 count: usize,
387 ) -> Poll<IoResult<usize>> {
388 let this = self.get_mut();
389 if this.state != WriterState::Open {
390 return Poll::Ready(Err(this.closed_io_error()));
391 }
392 if let Err(error) = this.check_write_limit(count) {
393 return Poll::Ready(Err(error));
394 }
395 match unsafe { this.session.as_mut().poll_write_unchecked(cx, input, index, count) } {
398 Poll::Ready(Ok(written)) => {
399 if let Err(error) = this.record_written_bytes(written) {
400 this.state = WriterState::Indeterminate;
401 return Poll::Ready(Err(error));
402 }
403 Poll::Ready(Ok(written))
404 }
405 Poll::Ready(Err(error)) => {
406 this.state = WriterState::Indeterminate;
407 Poll::Ready(Err(error))
408 }
409 Poll::Pending => Poll::Pending,
410 }
411 }
412
413 fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<IoResult<()>> {
414 let this = self.get_mut();
415 if this.state != WriterState::Open {
416 return Poll::Ready(Err(this.closed_io_error()));
417 }
418 match this.session.as_mut().poll_flush(cx) {
419 Poll::Ready(Ok(())) => Poll::Ready(Ok(())),
420 Poll::Ready(Err(error)) => {
421 this.state = WriterState::Indeterminate;
422 Poll::Ready(Err(error))
423 }
424 Poll::Pending => Poll::Pending,
425 }
426 }
427}
428
429impl Debug for AsyncFileWriter {
430 #[inline]
431 fn fmt(&self, formatter: &mut Formatter<'_>) -> FmtResult {
432 formatter
433 .debug_struct("AsyncFileWriter")
434 .field("info", &self.info)
435 .field("state", &self.state)
436 .finish_non_exhaustive()
437 }
438}
439
440impl Drop for AsyncFileWriter {
441 fn drop(&mut self) {
442 if !self.abort_completed
443 && matches!(
444 self.state,
445 WriterState::Open | WriterState::NotPublished | WriterState::Published
446 )
447 {
448 self.session.as_mut().cancel_on_drop();
449 }
450 }
451}