A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
a11::stores::ChunkStoreReader::State Struct Referenceabstract
Inheritance diagram for a11::stores::ChunkStoreReader::State:
[legend]

Classes

struct  Completion
 
struct  Request
 

Public Types

enum class  CompletionKind { kFragment , kEnd , kError }
 
enum class  Operation { kNone , kFetch , kClear }
 
using NextResult = std::optional< data::NodeFragment >
 

Public Member Functions

 State (std::shared_ptr< ChunkStore > chunk_store, ChunkStoreReaderOptions reader_options)
 
void Wake ()
 
void Cancel ()
 
std::vector< data::NodeFragment > TakeBuffered (size_t limit)
 Pop up to limit already-prefetched fragments, without ever waiting.
 
a11::Future< NextResult > Next (absl::Duration timeout)
 
void Drive ()
 Pump until there is nothing left to do without waiting.
 
bool DriveOnce ()
 
void InstallFetch (const a11::Future< data::NodeFragment > &pending, std::uint64_t position_wanted, std::uint64_t generation)
 
bool FetchArrived (std::uint64_t arrived_position, std::uint64_t generation, absl::StatusOr< data::NodeFragment > result, bool inline_drive=false)
 Record a completed fetch, in order.
 
void InstallClear (const a11::Future< data::NodeFragment > &pending, std::uint64_t generation)
 
void ClearDone (std::uint64_t generation, const absl::StatusOr< data::NodeFragment > &result)
 
void MaybeCompleteDone ()
 
void DrainArrivedLocked (std::vector< Completion > *completions, bool *start_clear, std::uint32_t *clear_seq, std::uint64_t *clear_generation) ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu)
 Apply arrived results that are now at the head, in position order.
 
void FinishFragmentLocked (data::NodeFragment fragment, std::vector< Completion > *completions) ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu)
 
size_t ActivePendingReadCountLocked () const ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu)
 
bool HasReadCapacityLocked (size_t about_to_issue=0) const ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu)
 
std::shared_ptr< Request > PopPendingReadLocked () ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu)
 
void CollectAvailableLocked (std::vector< Completion > *completions) ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu)
 
std::uint64_t position ABSL_GUARDED_BY (mu)
 
std::uint64_t chunks_read ABSL_GUARDED_BY (mu)=0
 
std::string current_mimetype ABSL_GUARDED_BY (mu)
 
std::optional< absl::Status > status ABSL_GUARDED_BY (mu)
 
std::deque< data::NodeFragment > buffer ABSL_GUARDED_BY (mu)
 
std::deque< std::shared_ptr< Request > > pending_reads ABSL_GUARDED_BY (mu)
 
bool queued ABSL_GUARDED_BY (mu)
 
Operation operation ABSL_GUARDED_BY (mu)
 
std::uint64_t operation_generation ABSL_GUARDED_BY (mu)=0
 
a11::Future< data::NodeFragment > active_operation ABSL_GUARDED_BY (mu)
 
std::uint64_t next_fetch_position ABSL_GUARDED_BY (mu)
 
size_t fetches_in_flight ABSL_GUARDED_BY (mu)=0
 
std::uint64_t fetch_generation ABSL_GUARDED_BY (mu)=0
 
std::map< std::uint64_t, absl::StatusOr< data::NodeFragment > > arrived ABSL_GUARDED_BY (mu)
 
std::vector< a11::Future< data::NodeFragment > > active_fetches ABSL_GUARDED_BY (mu)
 
bool done_completed ABSL_GUARDED_BY (mu)
 

Static Public Member Functions

static a11::internal::CallbackScheduler & Scheduler ()
 
static void Complete (std::vector< Completion > completions)
 

Public Attributes

const std::shared_ptr< ChunkStore > store
 
const ChunkStoreReaderOptions options
 
thread::Mutex mu
 
InlinePumpState pump
 
a11::Promise< a11::Unit > done_promise
 
const a11::Task done
 

Static Public Attributes

static constexpr size_t kMaxFetchesInFlight = 16
 How many reads the pump keeps outstanding at once.
 

Member Typedef Documentation

◆ NextResult

Member Enumeration Documentation

◆ CompletionKind

Enumerator
kFragment 
kEnd 
kError 

◆ Operation

Enumerator
kNone 
kFetch 
kClear 

Constructor & Destructor Documentation

◆ State()

a11::stores::ChunkStoreReader::State::State ( std::shared_ptr< ChunkStore >  chunk_store,
ChunkStoreReaderOptions  reader_options 
)
inline

Member Function Documentation

◆ ABSL_GUARDED_BY() [1/16]

std::uint64_t position a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [2/16]

std::string current_mimetype a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [3/16]

std::optional< absl::Status > status a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [4/16]

std::deque< data::NodeFragment > buffer a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [5/16]

std::deque< std::shared_ptr< Request > > pending_reads a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [6/16]

bool queued a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [7/16]

Operation operation a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [8/16]

a11::Future< data::NodeFragment > active_operation a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [9/16]

std::uint64_t next_fetch_position a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [10/16]

std::map< std::uint64_t, absl::StatusOr< data::NodeFragment > > arrived a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [11/16]

std::vector< a11::Future< data::NodeFragment > > active_fetches a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [12/16]

bool done_completed a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [13/16]

std::uint64_t chunks_read a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )
pure virtual

◆ ABSL_GUARDED_BY() [14/16]

std::uint64_t operation_generation a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )
pure virtual

◆ ABSL_GUARDED_BY() [15/16]

size_t fetches_in_flight a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )
pure virtual

◆ ABSL_GUARDED_BY() [16/16]

std::uint64_t fetch_generation a11::stores::ChunkStoreReader::State::ABSL_GUARDED_BY ( mu  )
pure virtual

◆ ActivePendingReadCountLocked()

size_t a11::stores::ChunkStoreReader::State::ActivePendingReadCountLocked ( ) const
inline

◆ Cancel()

void a11::stores::ChunkStoreReader::State::Cancel ( )
inline

◆ ClearDone()

void a11::stores::ChunkStoreReader::State::ClearDone ( std::uint64_t  generation,
const absl::StatusOr< data::NodeFragment > &  result 
)
inline

◆ CollectAvailableLocked()

void a11::stores::ChunkStoreReader::State::CollectAvailableLocked ( std::vector< Completion > *  completions)
inline

◆ Complete()

static void a11::stores::ChunkStoreReader::State::Complete ( std::vector< Completion >  completions)
inlinestatic

◆ DrainArrivedLocked()

void a11::stores::ChunkStoreReader::State::DrainArrivedLocked ( std::vector< Completion > *  completions,
bool *  start_clear,
std::uint32_t *  clear_seq,
std::uint64_t *  clear_generation 
)
inline

Apply arrived results that are now at the head, in position order.

The decision for each one is exactly what it was when fetches were serialised – end of stream, error, clear-then-buffer, or buffer – only now it is taken when the result reaches the front of the queue rather than when it happens to come back. Stops at the first gap, and stops as soon as a terminal status is set, leaving anything later in arrived to be discarded.

◆ Drive()

void a11::stores::ChunkStoreReader::State::Drive ( )
inline

Pump until there is nothing left to do without waiting.

A store that already holds the fragment resolves Get() inline, which is the normal case, so this keeps going for as long as that holds. Handing each such fragment back to the scheduler instead would cost a scheduler hop per fragment and leave the buffer no more than one fragment ahead, which is what would make NextMany/next_fragments return batches of one.

Uses a loop because callbacks run on pooled fibers with small stacks; recursion per fragment could overflow them.

◆ DriveOnce()

bool a11::stores::ChunkStoreReader::State::DriveOnce ( )
inline
Returns
Whether the caller should drive again without yielding.

◆ FetchArrived()

bool a11::stores::ChunkStoreReader::State::FetchArrived ( std::uint64_t  arrived_position,
std::uint64_t  generation,
absl::StatusOr< data::NodeFragment >  result,
bool  inline_drive = false 
)
inline

Record a completed fetch, in order.

Completed reads wait in arrived until their position is next. Errors follow the same ordering, so a later NotFound cannot end the stream while an earlier read remains outstanding.

Parameters
arrived_positionArrival-order position of this fetch.
generationFetch generation used to reject stale completions.
resultCompleted fragment or fetch error.
inline_driveTrue when the fetch resolved without waiting and DriveOnce() is still on the stack; then this reports back whether to keep driving instead of waking the scheduler.
Returns
Whether the caller should drive again immediately.

◆ FinishFragmentLocked()

void a11::stores::ChunkStoreReader::State::FinishFragmentLocked ( data::NodeFragment  fragment,
std::vector< Completion > *  completions 
)
inline

◆ HasReadCapacityLocked()

bool a11::stores::ChunkStoreReader::State::HasReadCapacityLocked ( size_t  about_to_issue = 0) const
inline

◆ InstallClear()

void a11::stores::ChunkStoreReader::State::InstallClear ( const a11::Future< data::NodeFragment > &  pending,
std::uint64_t  generation 
)
inline

◆ InstallFetch()

void a11::stores::ChunkStoreReader::State::InstallFetch ( const a11::Future< data::NodeFragment > &  pending,
std::uint64_t  position_wanted,
std::uint64_t  generation 
)
inline

◆ MaybeCompleteDone()

void a11::stores::ChunkStoreReader::State::MaybeCompleteDone ( )
inline

◆ Next()

a11::Future< NextResult > a11::stores::ChunkStoreReader::State::Next ( absl::Duration  timeout)
inline

◆ PopPendingReadLocked()

std::shared_ptr< Request > a11::stores::ChunkStoreReader::State::PopPendingReadLocked ( )
inline

◆ Scheduler()

static a11::internal::CallbackScheduler & a11::stores::ChunkStoreReader::State::Scheduler ( )
inlinestatic

◆ TakeBuffered()

std::vector< data::NodeFragment > a11::stores::ChunkStoreReader::State::TakeBuffered ( size_t  limit)
inline

Pop up to limit already-prefetched fragments, without ever waiting.

Only takes anything when nobody is queued ahead: pending_reads is served strictly in order, and a batch read that jumped that queue would reorder two concurrent readers. Everything it returns is a fragment CollectAvailableLocked() would have handed out next, in the same order, so this adds no semantics of its own. Terminal state (end of stream or error) is not handled here; NextMany falls back to the single-fragment path for that, so there is exactly one place that decides what the end of a stream looks like.

◆ Wake()

void a11::stores::ChunkStoreReader::State::Wake ( )
inline

Member Data Documentation

◆ done

const a11::Task a11::stores::ChunkStoreReader::State::done

◆ done_promise

a11::Promise<a11::Unit> a11::stores::ChunkStoreReader::State::done_promise

◆ kMaxFetchesInFlight

constexpr size_t a11::stores::ChunkStoreReader::State::kMaxFetchesInFlight = 16
staticconstexpr

How many reads the pump keeps outstanding at once.

With one, a store whose Get() has any latency delivers at one fragment per round trip and the prefetch buffer never fills; the reader spends its life waiting. Beyond about this many the extra concurrency stops buying anything and starts costing memory in the reorder map.

◆ mu

thread::Mutex a11::stores::ChunkStoreReader::State::mu
mutable

◆ options

const ChunkStoreReaderOptions a11::stores::ChunkStoreReader::State::options

◆ pump

InlinePumpState a11::stores::ChunkStoreReader::State::pump

◆ store

const std::shared_ptr<ChunkStore> a11::stores::ChunkStoreReader::State::store

The documentation for this struct was generated from the following file: