|
| | 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) |
| |
| 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.
| 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.
| 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.