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

Classes

struct  ReadCompletion
 
struct  ReadRequest
 

Public Types

enum class  ReadKind { kSequence , kArrivalOrder }
 

Public Member Functions

 State (std::string id)
 
a11::Future< data::NodeFragment > Read (ReadKind kind, std::uint64_t value, absl::Time deadline)
 
std::optional< absl::StatusOr< data::NodeFragment > > LookupReadLocked (const ReadRequest &request) const ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu)
 
void CollectReadsLocked (std::vector< ReadCompletion > *completions) ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu)
 
void RemoveRead (const ReadRequest *absl_nonnull request)
 
std::optional< absl::Status > CollectNext (std::vector< std::optional< data::NodeFragment > > &fragments, size_t limit, std::shared_ptr< thread::PermanentEvent > *waiter) ABSL_LOCKS_EXCLUDED(mu)
 Append whatever Next can return right now, without ever blocking.
 
absl::flat_hash_map< std::uint32_t, std::uint64_t > seq_to_arrival_order ABSL_GUARDED_BY (mu)
 
absl::flat_hash_map< std::uint64_t, std::uint32_t > arrival_order_to_seq ABSL_GUARDED_BY (mu)
 
absl::flat_hash_map< std::uint32_t, data::Chunk > chunks ABSL_GUARDED_BY (mu)
 
std::optional< std::uint32_t > final_seq ABSL_GUARDED_BY (mu)
 
std::uint64_t total_chunks_put ABSL_GUARDED_BY (mu)=0
 
std::uint64_t total_chunks_read ABSL_GUARDED_BY (mu)=0
 
std::optional< absl::Status > status ABSL_GUARDED_BY (mu)
 
std::vector< std::shared_ptr< ReadRequest > > pending_reads ABSL_GUARDED_BY (mu)
 
std::shared_ptr< thread::PermanentEvent > changed ABSL_GUARDED_BY (mu)
 

Static Public Member Functions

static void CompleteReads (std::vector< ReadCompletion > completions)
 

Public Attributes

thread::Mutex mu
 
const std::string node_id
 

Member Enumeration Documentation

◆ ReadKind

Enumerator
kSequence 
kArrivalOrder 

Constructor & Destructor Documentation

◆ State()

a11::stores::LocalChunkStore::State::State ( std::string  id)
inlineexplicit

Member Function Documentation

◆ ABSL_GUARDED_BY() [1/9]

absl::flat_hash_map< std::uint32_t, std::uint64_t > seq_to_arrival_order a11::stores::LocalChunkStore::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [2/9]

absl::flat_hash_map< std::uint64_t, std::uint32_t > arrival_order_to_seq a11::stores::LocalChunkStore::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [3/9]

absl::flat_hash_map< std::uint32_t, data::Chunk > chunks a11::stores::LocalChunkStore::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [4/9]

std::optional< std::uint32_t > final_seq a11::stores::LocalChunkStore::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [5/9]

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

◆ ABSL_GUARDED_BY() [6/9]

std::vector< std::shared_ptr< ReadRequest > > pending_reads a11::stores::LocalChunkStore::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [7/9]

std::shared_ptr< thread::PermanentEvent > changed a11::stores::LocalChunkStore::State::ABSL_GUARDED_BY ( mu  )

◆ ABSL_GUARDED_BY() [8/9]

std::uint64_t total_chunks_put a11::stores::LocalChunkStore::State::ABSL_GUARDED_BY ( mu  )
pure virtual

◆ ABSL_GUARDED_BY() [9/9]

std::uint64_t total_chunks_read a11::stores::LocalChunkStore::State::ABSL_GUARDED_BY ( mu  )
pure virtual

◆ CollectNext()

std::optional< absl::Status > a11::stores::LocalChunkStore::State::CollectNext ( std::vector< std::optional< data::NodeFragment > > &  fragments,
size_t  limit,
std::shared_ptr< thread::PermanentEvent > *  waiter 
)
inline

Append whatever Next can return right now, without ever blocking.

Split out of Next() so the common case – the fragments are already here – can run on the caller's thread. Wrapping the whole loop in a Submit() would spawn a fiber even when nothing needs waiting for, and a fiber spawn costs more than an order of magnitude what the read itself does.

Parameters
fragmentsAccumulator, appended to. Carries across waits, so a partial batch collected before a wait is still there after it.
limitMaximum fragments to accumulate, counting what is already there.
waiterSet to the generation event only when the caller must wait. It is snapshotted inside the lock, which is what closes the lost-wakeup window – see WaitForChange().
Returns
An engaged status when the call has its answer: OK means fragments is that answer, and a non-OK status is the terminal one to fail with. std::nullopt means the caller must wait on *waiter and try again.

◆ CollectReadsLocked()

void a11::stores::LocalChunkStore::State::CollectReadsLocked ( std::vector< ReadCompletion > *  completions)
inline

◆ CompleteReads()

static void a11::stores::LocalChunkStore::State::CompleteReads ( std::vector< ReadCompletion >  completions)
inlinestatic

◆ LookupReadLocked()

std::optional< absl::StatusOr< data::NodeFragment > > a11::stores::LocalChunkStore::State::LookupReadLocked ( const ReadRequest &  request) const
inline

◆ Read()

a11::Future< data::NodeFragment > a11::stores::LocalChunkStore::State::Read ( ReadKind  kind,
std::uint64_t  value,
absl::Time  deadline 
)
inline

◆ RemoveRead()

void a11::stores::LocalChunkStore::State::RemoveRead ( const ReadRequest *absl_nonnull  request)
inline

Member Data Documentation

◆ mu

thread::Mutex a11::stores::LocalChunkStore::State::mu

◆ node_id

const std::string a11::stores::LocalChunkStore::State::node_id

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