pub struct InMemoryGroupEngine { /* private fields */ }Implementations§
Source§impl InMemoryGroupEngine
impl InMemoryGroupEngine
pub fn with_cold_store(cold_store: ColdStoreHandle) -> Self
pub fn cold_store(&self) -> Option<ColdStoreHandle>
pub fn apply_committed_write( &mut self, command: GroupWriteCommand, placement: ShardPlacement, ) -> Result<GroupWriteResponse, GroupEngineError>
Sourcepub fn apply_stream_command(
&mut self,
command: StreamCommand,
placement: ShardPlacement,
) -> Result<GroupWriteResponse, GroupEngineError>
pub fn apply_stream_command( &mut self, command: StreamCommand, placement: ShardPlacement, ) -> Result<GroupWriteResponse, GroupEngineError>
Applies one canonical StreamCommand to the deterministic state
machine and lifts its StreamResponse into the group-level response,
maintaining the group commit index and per-stream append counts.
pub fn check_cold_write_admission_bytes( &self, stream_id: &BucketStreamId, admission: ColdWriteAdmission, incoming_bytes: u64, ) -> Result<(), GroupEngineError>
pub fn access_requires_write( &self, stream_id: &BucketStreamId, now_ms: u64, renew_ttl: bool, ) -> Result<bool, GroupEngineError>
pub fn read_stream_plan( &mut self, request: &ReadStreamRequest, placement: ShardPlacement, ) -> Result<StreamReadPlan, GroupEngineError>
pub fn read_stream_plan_after_access( &self, request: &ReadStreamRequest, ) -> Result<StreamReadPlan, GroupEngineError>
pub fn head_stream_after_access( &mut self, request: &HeadStreamRequest, placement: ShardPlacement, ) -> Result<HeadStreamResponse, GroupEngineError>
pub fn get_stream_attrs_after_access( &mut self, request: &GetStreamAttrsRequest, placement: ShardPlacement, ) -> Result<GetStreamAttrsResponse, GroupEngineError>
pub async fn read_payload_from_plan( cold_store: Option<&ColdStoreHandle>, cold_index_cache: Option<&Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>, stream_id: &BucketStreamId, plan: &StreamReadPlan, ) -> Result<Vec<u8>, GroupEngineError>
pub fn stream_tail_offset(&self, stream_id: &BucketStreamId) -> Option<u64>
Trait Implementations§
Source§impl Clone for InMemoryGroupEngine
impl Clone for InMemoryGroupEngine
Source§fn clone(&self) -> InMemoryGroupEngine
fn clone(&self) -> InMemoryGroupEngine
Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
Performs copy-assignment from
source. Read moreSource§impl Debug for InMemoryGroupEngine
impl Debug for InMemoryGroupEngine
Source§impl Default for InMemoryGroupEngine
impl Default for InMemoryGroupEngine
Source§fn default() -> InMemoryGroupEngine
fn default() -> InMemoryGroupEngine
Returns the “default value” for a type. Read more
Source§impl GroupEngine for InMemoryGroupEngine
impl GroupEngine for InMemoryGroupEngine
fn create_stream<'a>( &'a mut self, request: CreateStreamRequest, placement: ShardPlacement, admission: ColdWriteAdmission, ) -> GroupCreateStreamFuture<'a>
fn create_stream_external<'a>( &'a mut self, request: CreateStreamExternalRequest, placement: ShardPlacement, ) -> GroupCreateStreamFuture<'a>
fn read_stream<'a>( &'a mut self, request: ReadStreamRequest, placement: ShardPlacement, ) -> GroupReadStreamFuture<'a>
fn read_stream_parts<'a>( &'a mut self, request: ReadStreamRequest, placement: ShardPlacement, ) -> GroupReadStreamPartsFuture<'a>
fn publish_snapshot<'a>( &'a mut self, request: PublishSnapshotRequest, placement: ShardPlacement, ) -> GroupPublishSnapshotFuture<'a>
fn read_snapshot<'a>( &'a mut self, request: ReadSnapshotRequest, placement: ShardPlacement, ) -> GroupReadSnapshotFuture<'a>
fn delete_snapshot<'a>( &'a mut self, request: DeleteSnapshotRequest, placement: ShardPlacement, ) -> GroupDeleteSnapshotFuture<'a>
fn bootstrap_stream<'a>( &'a mut self, request: BootstrapStreamRequest, placement: ShardPlacement, ) -> GroupBootstrapStreamFuture<'a>
fn touch_stream_access<'a>( &'a mut self, stream_id: BucketStreamId, now_ms: u64, renew_ttl: bool, placement: ShardPlacement, ) -> GroupTouchStreamAccessFuture<'a>
fn head_stream<'a>( &'a mut self, request: HeadStreamRequest, placement: ShardPlacement, ) -> GroupHeadStreamFuture<'a>
fn get_stream_attrs<'a>( &'a mut self, request: GetStreamAttrsRequest, placement: ShardPlacement, ) -> GroupGetStreamAttrsFuture<'a>
fn update_stream_attrs<'a>( &'a mut self, request: UpdateStreamAttrsRequest, placement: ShardPlacement, ) -> GroupUpdateStreamAttrsFuture<'a>
fn close_stream<'a>( &'a mut self, request: CloseStreamRequest, placement: ShardPlacement, ) -> GroupCloseStreamFuture<'a>
fn delete_stream<'a>( &'a mut self, request: DeleteStreamRequest, placement: ShardPlacement, ) -> GroupDeleteStreamFuture<'a>
Source§fn ack_cold_gc<'a>(
&'a mut self,
up_to_seq: u64,
placement: ShardPlacement,
) -> GroupAckColdGcFuture<'a>
fn ack_cold_gc<'a>( &'a mut self, up_to_seq: u64, placement: ShardPlacement, ) -> GroupAckColdGcFuture<'a>
Replicated confirmation that cold-GC entries up to
up_to_seq have been
physically reclaimed; pops them from the queue. Default unsupported.Source§fn plan_cold_gc<'a>(
&'a mut self,
max: usize,
_placement: ShardPlacement,
) -> GroupPlanColdGcFuture<'a>
fn plan_cold_gc<'a>( &'a mut self, max: usize, _placement: ShardPlacement, ) -> GroupPlanColdGcFuture<'a>
Leader-local read of the front of the cold-GC queue for the background
worker to reclaim. Default returns an empty batch.
fn append<'a>( &'a mut self, request: AppendRequest, placement: ShardPlacement, admission: ColdWriteAdmission, ) -> GroupAppendFuture<'a>
fn append_external<'a>( &'a mut self, request: AppendExternalRequest, placement: ShardPlacement, ) -> GroupAppendFuture<'a>
fn append_batch<'a>( &'a mut self, request: AppendBatchRequest, placement: ShardPlacement, admission: ColdWriteAdmission, ) -> GroupAppendBatchFuture<'a>
fn flush_cold<'a>( &'a mut self, request: FlushColdRequest, placement: ShardPlacement, ) -> GroupFlushColdFuture<'a>
fn plan_cold_flush<'a>( &'a mut self, request: PlanColdFlushRequest, _placement: ShardPlacement, ) -> GroupPlanColdFlushFuture<'a>
fn plan_next_cold_flush_batch<'a>( &'a mut self, request: PlanGroupColdFlushRequest, _placement: ShardPlacement, max_candidates: usize, ) -> GroupPlanNextColdFlushBatchFuture<'a>
fn cold_hot_backlog<'a>( &'a mut self, stream_id: BucketStreamId, _placement: ShardPlacement, ) -> GroupColdHotBacklogFuture<'a>
fn snapshot<'a>( &'a mut self, placement: ShardPlacement, ) -> GroupSnapshotFuture<'a>
fn install_snapshot<'a>( &'a mut self, snapshot: GroupSnapshot, ) -> GroupInstallSnapshotFuture<'a>
fn accepts_local_writes(&self) -> bool
fn require_local_live_read_owner<'a>( &'a mut self, _placement: ShardPlacement, ) -> GroupRequireLiveReadOwnerFuture<'a>
fn append_batch_many<'a>( &'a mut self, requests: Vec<AppendBatchRequest>, placement: ShardPlacement, admission: ColdWriteAdmission, ) -> GroupWriteBatchFuture<'a>
fn shutdown<'a>(&'a mut self) -> GroupShutdownFuture<'a>
fn write_batch<'a>( &'a mut self, commands: Vec<GroupWriteCommand>, placement: ShardPlacement, ) -> GroupWriteBatchFuture<'a>
Source§fn dispatch_stream_command<'a>(
&'a mut self,
command: StreamCommand,
placement: ShardPlacement,
) -> GroupWriteFuture<'a>
fn dispatch_stream_command<'a>( &'a mut self, command: StreamCommand, placement: ShardPlacement, ) -> GroupWriteFuture<'a>
Routes one canonical
StreamCommand to the matching typed engine
method. This is the only spelling of the command-to-operation mapping;
engines that replicate commands wholesale (raft) bypass it.Auto Trait Implementations§
impl !RefUnwindSafe for InMemoryGroupEngine
impl !UnwindSafe for InMemoryGroupEngine
impl Freeze for InMemoryGroupEngine
impl Send for InMemoryGroupEngine
impl Sync for InMemoryGroupEngine
impl Unpin for InMemoryGroupEngine
impl UnsafeUnpin for InMemoryGroupEngine
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
Mutably borrows from an owned value. Read more