1use std::fmt::Debug;
11use std::fmt::Formatter;
12use std::fmt::Result as FmtResult;
13use std::io::Error as IoError;
14use std::io::ErrorKind as IoErrorKind;
15use std::io::Result as IoResult;
16
17use qubit_io::Output;
18
19use crate::error::FsEffectState;
20use crate::error::FsError;
21use crate::error::FsErrorKind;
22use crate::error::FsOperation;
23use crate::error::FsResult;
24use crate::facade::facade_core::FacadeCore;
25use crate::facade::internal::ByteBudget;
26use crate::facade::internal::FileSystemResource;
27use crate::metadata::AchievedAtomicity;
28use crate::metadata::AtomicityRequirement;
29use crate::metadata::DurabilityRequirement;
30use crate::metadata::OpenedFileInfo;
31use crate::metadata::WriteOutcome;
32use crate::spi::FileWriterSpi;
33use crate::write::WriteAbortOutcome;
34use crate::write::WriteFailure;
35use crate::write::WriteFailureState;
36use crate::write::WriterState;
37
38pub struct FileWriter {
60 session: Box<dyn FileWriterSpi>,
62 info: OpenedFileInfo,
64 state: WriterState,
66 abort_completed: bool,
68 atomicity: AtomicityRequirement,
70 durability: DurabilityRequirement,
72 provider: Box<str>,
74 write_budget: Option<ByteBudget>,
76 written_bytes: u64,
78}
79
80impl FileWriter {
81 #[inline]
90 #[must_use]
91 pub(crate) fn new(
92 info: OpenedFileInfo,
93 session: Box<dyn FileWriterSpi>,
94 atomicity: AtomicityRequirement,
95 durability: DurabilityRequirement,
96 provider: &str,
97 max_write_bytes: Option<u64>,
98 ) -> Self {
99 Self {
100 session,
101 info,
102 state: WriterState::Open,
103 abort_completed: false,
104 atomicity,
105 durability,
106 provider: provider.into(),
107 write_budget: max_write_bytes
108 .map(|maximum| FacadeCore::byte_budget(FileSystemResource::WriteBytes, maximum)),
109 written_bytes: 0,
110 }
111 }
112
113 #[inline]
118 #[must_use]
119 pub fn info(&self) -> &OpenedFileInfo {
120 &self.info
121 }
122
123 #[inline]
128 #[must_use]
129 pub const fn state(&self) -> WriterState {
130 self.state
131 }
132
133 #[inline]
135 #[must_use]
136 pub(crate) const fn written_bytes(&self) -> u64 {
137 self.written_bytes
138 }
139
140 pub fn commit(&mut self) -> Result<WriteOutcome, WriteFailure> {
153 if self.state != WriterState::Open {
154 let publication_state = self.state.publication_failure_state();
155 return Err(WriteFailure::new(
156 self.invalid_state(
157 FsOperation::CommitWriter,
158 "writer cannot be committed in its current state",
159 ),
160 publication_state,
161 ));
162 }
163 let outcome = self.session.commit();
164 match outcome {
165 Ok(outcome) => {
166 self.state = WriterState::Committed;
167 if self.atomicity == AtomicityRequirement::Required && outcome.atomicity() != AchievedAtomicity::Atomic
168 {
169 self.state = WriterState::Published;
170 return Err(WriteFailure::new(
171 FsError::new(
172 FsErrorKind::ProviderContractViolation,
173 FsOperation::CommitWriter,
174 "provider reported non-atomic success for an atomic-required write",
175 )
176 .with_path(self.info.path().clone())
177 .with_provider(&self.provider)
178 .with_effect_state(FsEffectState::Applied),
179 WriteFailureState::Published,
180 ));
181 }
182 if self.durability == DurabilityRequirement::Required && !outcome.durable() {
183 self.state = WriterState::Published;
184 return Err(WriteFailure::new(
185 FsError::new(
186 FsErrorKind::ProviderContractViolation,
187 FsOperation::CommitWriter,
188 "provider reported non-durable success for a durability-required write",
189 )
190 .with_path(self.info.path().clone())
191 .with_provider(&self.provider)
192 .with_effect_state(FsEffectState::Applied),
193 WriteFailureState::Published,
194 ));
195 }
196 if let Some(bytes_written) = outcome.bytes_written()
197 && bytes_written != self.written_bytes
198 {
199 self.state = WriterState::Published;
200 return Err(WriteFailure::new(
201 FsError::new(
202 FsErrorKind::ProviderContractViolation,
203 FsOperation::CommitWriter,
204 "provider reported a byte count different from the bytes accepted by the writer",
205 )
206 .with_path(self.info.path().clone())
207 .with_provider(&self.provider)
208 .with_effect_state(FsEffectState::Applied),
209 WriteFailureState::Published,
210 ));
211 }
212 Ok(outcome)
213 }
214 Err(failure) => {
215 self.state = match failure.state() {
216 WriteFailureState::RetryableNotPublished => WriterState::Open,
217 WriteFailureState::NotPublished => WriterState::NotPublished,
218 WriteFailureState::Published => WriterState::Published,
219 WriteFailureState::Indeterminate => WriterState::Indeterminate,
220 };
221 let (error, state) = failure.into_parts();
222 Err(WriteFailure::new(
223 self.contextual_error(error, FsOperation::CommitWriter),
224 state,
225 ))
226 }
227 }
228 }
229
230 pub fn abort(&mut self) -> FsResult<WriteAbortOutcome> {
244 if self.abort_completed
245 || !matches!(
246 self.state,
247 WriterState::Open | WriterState::NotPublished | WriterState::Published | WriterState::Indeterminate
248 )
249 {
250 return Err(self.invalid_state(
251 FsOperation::AbortWriter,
252 "writer cannot be aborted in its current state",
253 ));
254 }
255 match self.session.abort() {
256 Ok(outcome) => {
257 self.abort_completed = true;
258 self.state = match outcome {
259 WriteAbortOutcome::NotPublished => WriterState::Aborted,
260 WriteAbortOutcome::Published => WriterState::Published,
261 WriteAbortOutcome::Indeterminate => WriterState::Indeterminate,
262 };
263 Ok(outcome)
264 }
265 Err(error) => {
266 if error.has_indeterminate_effect() {
267 self.state = WriterState::Indeterminate;
268 }
269 Err(self.contextual_error(error, FsOperation::AbortWriter))
270 }
271 }
272 }
273
274 fn invalid_state(&self, operation: FsOperation, message: &str) -> FsError {
276 FsError::new(FsErrorKind::InvalidState, operation, message)
277 .with_path(self.info.path().clone())
278 .with_provider(&self.provider)
279 }
280
281 fn closed_io_error(&self) -> IoError {
283 IoError::new(
284 IoErrorKind::BrokenPipe,
285 self.invalid_state(FsOperation::Write, "writer no longer accepts bytes"),
286 )
287 }
288
289 fn check_write_limit(&self, count: usize) -> IoResult<u64> {
291 let count = FacadeCore::quantity_from_usize(count, FsOperation::Write, self.info.path(), &self.provider)
292 .map_err(FsError::into_io_error)?;
293 if let Some(budget) = &self.write_budget
294 && let Err(error) = budget.check_available(count)
295 {
296 return Err(FacadeCore::budget_error(
297 error,
298 FsOperation::Write,
299 self.info.path(),
300 &self.provider,
301 "write session exceeds the provider byte limit",
302 )
303 .into_io_error());
304 }
305 Ok(count)
306 }
307
308 fn record_written_bytes(&mut self, count: usize) -> IoResult<()> {
314 let count = FacadeCore::quantity_from_usize(count, FsOperation::Write, self.info.path(), &self.provider)
315 .map_err(FsError::into_io_error)?;
316 if let Some(error) = self
317 .write_budget
318 .as_mut()
319 .and_then(|budget| budget.try_consume(count).err())
320 {
321 return Err(FacadeCore::budget_error(
322 error,
323 FsOperation::Write,
324 self.info.path(),
325 &self.provider,
326 "write session exceeds the provider byte limit",
327 )
328 .into_io_error());
329 }
330 self.written_bytes = self
331 .written_bytes
332 .checked_add(count)
333 .ok_or_else(|| self.byte_count_error())?;
334 Ok(())
335 }
336
337 fn byte_count_error(&self) -> IoError {
340 FsError::new(
341 FsErrorKind::ResourceLimitExceeded,
342 FsOperation::Write,
343 "write byte count exceeds the filesystem API reporting range",
344 )
345 .with_path(self.info.path().clone())
346 .with_provider(&self.provider)
347 .into_io_error()
348 }
349
350 fn contextual_error(&self, error: FsError, operation: FsOperation) -> FsError {
352 error
353 .with_operation(operation)
354 .with_missing_context(self.info.path(), None, &self.provider)
355 }
356}
357
358impl Output for FileWriter {
359 type Item = u8;
360
361 #[inline]
362 fn is_buffered(&self) -> bool {
363 self.session.is_buffered()
364 }
365
366 unsafe fn write_unchecked(&mut self, input: &[u8], index: usize, count: usize) -> IoResult<usize> {
367 if self.state != WriterState::Open {
368 return Err(self.closed_io_error());
369 }
370 self.check_write_limit(count)?;
371 match unsafe { self.session.write_unchecked(input, index, count) } {
374 Ok(value) => {
375 if let Err(error) = self.record_written_bytes(value) {
376 self.state = WriterState::Indeterminate;
377 return Err(error);
378 }
379 Ok(value)
380 }
381 Err(error) => {
382 self.state = WriterState::Indeterminate;
383 Err(error)
384 }
385 }
386 }
387
388 fn flush(&mut self) -> IoResult<()> {
389 if self.state != WriterState::Open {
390 return Err(self.closed_io_error());
391 }
392 match self.session.flush() {
393 Ok(()) => Ok(()),
394 Err(error) => {
395 self.state = WriterState::Indeterminate;
396 Err(error)
397 }
398 }
399 }
400}
401
402impl Debug for FileWriter {
403 #[inline]
404 fn fmt(&self, formatter: &mut Formatter<'_>) -> FmtResult {
405 formatter
406 .debug_struct("FileWriter")
407 .field("info", &self.info)
408 .field("state", &self.state)
409 .finish_non_exhaustive()
410 }
411}
412
413impl Drop for FileWriter {
414 fn drop(&mut self) {
415 if !self.abort_completed
416 && matches!(
417 self.state,
418 WriterState::Open | WriterState::NotPublished | WriterState::Published
419 )
420 {
421 let _ = self.session.abort();
422 }
423 }
424}