pub struct Writer<B: Blob> { /* private fields */ }Expand description
Unique writer to a cache-wrapped Blob.
Implementations§
Source§impl<B: Blob> Writer<B>
impl<B: Blob> Writer<B>
Sourcepub async fn new(
blob: B,
original_blob_size: u64,
capacity: usize,
cache_ref: CacheRef,
) -> Result<Self, Error>
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.
Sourcepub async fn append(&mut self, buf: &[u8]) -> Result<u64, Error>
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.
Sourcepub async fn append_owned(&mut self, buf: IoBuf) -> Result<u64, Error>
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.
Sourcepub fn try_read_sync_into(&self, buf: &mut [u8], offset: u64) -> bool
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.
Sourcepub async fn read_at(&self, offset: u64, len: usize) -> Result<IoBufs, Error>
pub async fn read_at(&self, offset: u64, len: usize) -> Result<IoBufs, Error>
Read exactly len immutable bytes starting at offset.
Sourcepub async fn read_up_to(
&self,
offset: u64,
len: usize,
bufs: impl Into<IoBufMut> + Send,
) -> Result<(IoBufMut, usize), Error>
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.
Sourcepub async fn read_many_into(
&self,
buf: &mut [u8],
offsets: &[u64],
item_size: NonZeroUsize,
) -> Result<usize, Error>
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.
Sourcepub fn try_read_many_sync_into(
&self,
buf: &mut [u8],
offsets: &[u64],
item_size: NonZeroUsize,
) -> Vec<usize>
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.
Sourcepub fn try_read_ranges_sync_into(
&self,
buf: &mut [u8],
ranges: &[(u64, usize)],
) -> Vec<usize>
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.
Sourcepub async fn read_into(&self, buf: &mut [u8], offset: u64) -> Result<(), Error>
pub async fn read_into(&self, buf: &mut [u8], offset: u64) -> Result<(), Error>
Reads bytes starting at offset into buf.
Sourcepub async fn replay(
&mut self,
buffer_size: NonZeroUsize,
read_options: ReadOptions,
) -> Result<Replay<B>, Error>
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.
Sourcepub async fn snapshot(&mut self) -> Result<Sealed<B>, Error>
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.
Sourcepub async fn sync(&mut self) -> Result<(), Error>
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.
Sourcepub async fn start_sync(&mut self) -> Handle<()> ⓘ
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.
Sourcepub async fn recoverable_prefix_len(
&self,
proven: u64,
buffer_size: NonZeroUsize,
read_options: ReadOptions,
) -> Result<u64, Error>
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.
Sourcepub async fn read_range(
blob: &B,
logical_page_size: NonZeroU16,
offset: u64,
len: usize,
read_options: ReadOptions,
) -> Result<IoBufs, Error>
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.
Sourcepub async fn read_tail(
blob: &B,
physical_size: u64,
logical_page_size: NonZeroU16,
len: usize,
read_options: ReadOptions,
) -> Result<(u64, IoBufs), Error>
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.
Sourcepub async fn wait_for_sync(&mut self) -> Result<(), Error>
pub async fn wait_for_sync(&mut self) -> Result<(), Error>
Wait for any started sync to complete without starting a new sync.
Sourcepub async fn resize(&mut self, size: u64) -> Result<(), Error>
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.
Sourcepub async fn seal(self) -> Result<(Sealed<B>, Handle<()>), Error>
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> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
Source§fn with_current_context(self) -> WithContext<Self> ⓘ
fn with_current_context(self) -> WithContext<Self> ⓘ
impl<A, B, T> HttpServerConnExec<A, B> for Twhere
B: Body,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
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