Skip to main content

Writer

Struct Writer 

Source
pub struct Writer<B: Blob> { /* private fields */ }
Expand description

Unique writer to a cache-wrapped Blob.

Implementations§

Source§

impl<B: Blob> Writer<B>

Source

pub async fn new( blob: B, original_blob_size: u64, capacity: usize, cache_ref: CacheRef, ) -> Result<Self, Error>

Wrap blob in a Writer. blob must already hold original_blob_size physical bytes; reads are cached through cache_ref and appends stage in a write buffer of capacity capacity. Rewinds the blob if necessary so it only contains checksum-validated data.

The blob’s tail-page contents must be durable (freshly opened after a crash, or synced since the last partial-page rewrite): the discovered checksum slot seeds the writer’s durable-slot tracking, so wrapping a blob whose tail rewrite is still volatile would let a later unsynced flush overwrite the only durable slot.

Source

pub async fn append(&mut self, buf: &[u8]) -> Result<u64, Error>

Append all bytes in buf to the tip of the blob, returning the logical offset at which the first byte was written.

Source

pub async fn append_owned(&mut self, buf: IoBuf) -> Result<u64, Error>

Append owned bytes to the tip of the blob.

Large appends fill the current tip to a page boundary, write complete pages directly to the blob, and leave only a sub-page suffix in the write buffer. This avoids copying full-page payloads while preserving the invariant that the buffer starts at current_page.

Source

pub const fn size(&self) -> u64

Returns the size of the blob.

Source

pub fn try_read_sync_into(&self, buf: &mut [u8], offset: u64) -> bool

Read into buf if it can be done synchronously without I/O. Returns true only if all buf.len() bytes were satisfied from the page cache and/or the in-memory tail. When false is returned, the contents of buf are unspecified.

Source

pub async fn read_at(&self, offset: u64, len: usize) -> Result<IoBufs, Error>

Read exactly len immutable bytes starting at offset.

Source

pub async fn read_up_to( &self, offset: u64, len: usize, bufs: impl Into<IoBufMut> + Send, ) -> Result<(IoBufMut, usize), Error>

Reads up to len bytes starting at offset, but only as many as are available.

Returns the buffer (truncated to actual bytes read) and the number of bytes read. Returns an error if no bytes are available at the given offset.

Source

pub async fn read_many_into( &self, buf: &mut [u8], offsets: &[u64], item_size: NonZeroUsize, ) -> Result<usize, Error>

Read multiple fixed-size items at sorted byte offsets into a contiguous caller buffer.

buf must be exactly offsets.len() * item_size bytes. All offsets must be sorted, non-overlapping, and within bounds.

Returns the number of items fully served without a blob read (from the in-memory tail and the page cache). The remaining items required at least one blob read.

Source

pub fn try_read_many_sync_into( &self, buf: &mut [u8], offsets: &[u64], item_size: NonZeroUsize, ) -> Vec<usize>

Like Self::read_many_into, but synchronous and cache-only. Returns the indices of items that require a blob read. Their slots in buf hold unspecified bytes.

Source

pub fn try_read_ranges_sync_into( &self, buf: &mut [u8], ranges: &[(u64, usize)], ) -> Vec<usize>

Like Self::try_read_many_sync_into, but for variable-length (offset, len) ranges: buf holds one slot per range, back to back.

Source

pub async fn read_into(&self, buf: &mut [u8], offset: u64) -> Result<(), Error>

Reads bytes starting at offset into buf.

Source

pub async fn replay( &mut self, buffer_size: NonZeroUsize, read_options: ReadOptions, ) -> Result<Replay<B>, Error>

Flushes any buffered data, then returns a Replay for the underlying blob.

The returned replay can be used to sequentially read all pages from the blob while ensuring all data passes integrity verification. CRCs are validated but not included in the output. Every underlying blob read performed by the returned replay uses read_options, including refills after seeking.

This is not a durable operation. Buffered data may be plainly written so the replay can read it, but callers must still use sync if that data must survive a crash.

Source

pub async fn snapshot(&mut self) -> Result<Sealed<B>, Error>

Flush buffered data and capture an immutable super::Sealed view without consuming the writer.

This writes buffered bytes to the blob layout but does not make them durable. Call Self::sync if the returned handle’s bytes must survive a crash.

If this writer later rewinds or truncates into the returned handle’s range, reads from that handle may observe unspecified contents.

Source

pub async fn sync(&mut self) -> Result<(), Error>

Flushes buffered data and makes all pending mutations durable.

A newly flushed write can carry WriteOptions::SYNC when no earlier mutation is pending. Otherwise, Blob::sync provides the barrier for all pending mutations.

Source

pub async fn start_sync(&mut self) -> Handle<()>

Flushes buffered data and begins making all pending mutations durable, returning a completion handle.

Awaiting the returned Handle waits for the same durability guarantee as Self::sync for the state flushed by this call. Later calls to Self::sync and writer methods that mutate the blob first wait for any outstanding start_sync handles.

Source

pub async fn recoverable_prefix_len( &self, proven: u64, buffer_size: NonZeroUsize, read_options: ReadOptions, ) -> Result<u64, Error>

Length of the longest contiguous prefix of well-formed pages on the blob.

Self::new sizes a blob by scanning backward to its last valid page, which cannot detect an earlier page that was lost or corrupted. This scans forward instead, stopping at the first invalid or short page.

Expects all appended bytes to have reached the blob (as after recovery): a partial page still buffered in this writer is unreadable from the blob and fails the scan. buffer_size bounds each blob read, with a minimum of one physical page. Applies read_options to every blob read.

proven is a logical byte offset already known valid (a durability watermark or a replay-validated prefix). Pages wholly below it are accepted without reading, and the scan starts at the page containing it. A proof past the blob’s content clamps to the full pages that exist, so a partial tail is still read rather than credited as full.

Source

pub async fn read_range( blob: &B, logical_page_size: NonZeroU16, offset: u64, len: usize, read_options: ReadOptions, ) -> Result<IoBufs, Error>

Read a logical range directly from a raw paged blob, validating every page it spans.

Returns Error::BlobInsufficientLength when valid page contents do not cover the whole range, and Error::OffsetOverflow when its end or page offsets overflow.

Source

pub async fn read_tail( blob: &B, physical_size: u64, logical_page_size: NonZeroU16, len: usize, read_options: ReadOptions, ) -> Result<(u64, IoBufs), Error>

Read the terminal logical range of a raw paged blob and return its logical size.

The blob must contain only complete physical pages. The last page is validated first to determine the logical end. If len crosses a page boundary, preceding pages are validated with Self::read_range.

Source

pub async fn wait_for_sync(&mut self) -> Result<(), Error>

Wait for any started sync to complete without starting a new sync.

Source

pub async fn resize(&mut self, size: u64) -> Result<(), Error>

Resize the blob to the provided logical size.

This truncates the blob to contain only size logical bytes. The physical blob size will be adjusted to include the necessary CRC records for the remaining pages.

§Warning
  • Concurrent mutable operations (append, resize) are not supported and will cause data loss.
  • Concurrent readers which try to read past the new size during the resize may error.
  • The resize is not guaranteed durable until the next sync.
Source

pub async fn seal(self) -> Result<(Sealed<B>, Handle<()>), Error>

Consume the write handle, flushing buffered bytes and beginning a sync of the blob.

Returns an immutable super::Sealed read handle plus a completion handle for the started sync. Reads through the super::Sealed handle observe flushed bytes immediately; durability isn’t guaranteed until the sync handle completes.

Auto Trait Implementations§

§

impl<B> !RefUnwindSafe for Writer<B>

§

impl<B> !UnwindSafe for Writer<B>

§

impl<B> Freeze for Writer<B>
where B: Freeze,

§

impl<B> Send for Writer<B>

§

impl<B> Sync for Writer<B>

§

impl<B> Unpin for Writer<B>
where B: Unpin,

§

impl<B> UnsafeUnpin for Writer<B>
where B: UnsafeUnpin,

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
Source§

impl<A, B, T> HttpServerConnExec<A, B> for T
where B: Body,

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> Threaded<T> for T

Source§

type Rest = ()

The outputs beyond the threaded value.
Source§

fn split(self) -> (T, ())

Splits into the threaded value and the extra outputs.
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more