pub trait GroupEngine: Send + 'static {
Show 31 methods
// Required methods
fn create_stream<'a>(
&'a mut self,
request: CreateStreamRequest,
placement: ShardPlacement,
admission: ColdWriteAdmission,
) -> GroupCreateStreamFuture<'a>;
fn head_stream<'a>(
&'a mut self,
request: HeadStreamRequest,
placement: ShardPlacement,
) -> GroupHeadStreamFuture<'a>;
fn read_stream<'a>(
&'a mut self,
request: ReadStreamRequest,
placement: ShardPlacement,
) -> GroupReadStreamFuture<'a>;
fn touch_stream_access<'a>(
&'a mut self,
stream_id: BucketStreamId,
now_ms: u64,
renew_ttl: bool,
placement: ShardPlacement,
) -> GroupTouchStreamAccessFuture<'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>;
fn append<'a>(
&'a mut self,
request: AppendRequest,
placement: ShardPlacement,
admission: ColdWriteAdmission,
) -> GroupAppendFuture<'a>;
fn append_batch<'a>(
&'a mut self,
request: AppendBatchRequest,
placement: ShardPlacement,
admission: ColdWriteAdmission,
) -> GroupAppendBatchFuture<'a>;
fn snapshot<'a>(
&'a mut self,
placement: ShardPlacement,
) -> GroupSnapshotFuture<'a>;
fn install_snapshot<'a>(
&'a mut self,
snapshot: GroupSnapshot,
) -> GroupInstallSnapshotFuture<'a>;
// Provided methods
fn accepts_local_writes(&self) -> bool { ... }
fn create_stream_external<'a>(
&'a mut self,
request: CreateStreamExternalRequest,
_placement: ShardPlacement,
) -> GroupCreateStreamFuture<'a> { ... }
fn get_stream_attrs<'a>(
&'a mut self,
request: GetStreamAttrsRequest,
_placement: ShardPlacement,
) -> GroupGetStreamAttrsFuture<'a> { ... }
fn read_stream_parts<'a>(
&'a mut self,
request: ReadStreamRequest,
placement: ShardPlacement,
) -> GroupReadStreamPartsFuture<'a> { ... }
fn require_local_live_read_owner<'a>(
&'a mut self,
_placement: ShardPlacement,
) -> GroupRequireLiveReadOwnerFuture<'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 update_stream_attrs<'a>(
&'a mut self,
request: UpdateStreamAttrsRequest,
_placement: ShardPlacement,
) -> GroupUpdateStreamAttrsFuture<'a> { ... }
fn ack_cold_gc<'a>(
&'a mut self,
_up_to_seq: u64,
_placement: ShardPlacement,
) -> GroupAckColdGcFuture<'a> { ... }
fn plan_cold_gc<'a>(
&'a mut self,
_max: usize,
_placement: ShardPlacement,
) -> GroupPlanColdGcFuture<'a> { ... }
fn append_external<'a>(
&'a mut self,
request: AppendExternalRequest,
_placement: ShardPlacement,
) -> GroupAppendFuture<'a> { ... }
fn append_batch_many<'a>(
&'a mut self,
requests: Vec<AppendBatchRequest>,
placement: ShardPlacement,
admission: ColdWriteAdmission,
) -> GroupWriteBatchFuture<'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 shutdown<'a>(&'a mut self) -> GroupShutdownFuture<'a> { ... }
fn write_batch<'a>(
&'a mut self,
commands: Vec<GroupWriteCommand>,
placement: ShardPlacement,
) -> GroupWriteBatchFuture<'a> { ... }
fn dispatch_stream_command<'a>(
&'a mut self,
command: StreamCommand,
placement: ShardPlacement,
) -> GroupWriteFuture<'a> { ... }
}Required Methods§
fn create_stream<'a>( &'a mut self, request: CreateStreamRequest, placement: ShardPlacement, admission: ColdWriteAdmission, ) -> GroupCreateStreamFuture<'a>
fn head_stream<'a>( &'a mut self, request: HeadStreamRequest, placement: ShardPlacement, ) -> GroupHeadStreamFuture<'a>
fn read_stream<'a>( &'a mut self, request: ReadStreamRequest, placement: ShardPlacement, ) -> GroupReadStreamFuture<'a>
fn touch_stream_access<'a>( &'a mut self, stream_id: BucketStreamId, now_ms: u64, renew_ttl: bool, placement: ShardPlacement, ) -> GroupTouchStreamAccessFuture<'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>
fn append<'a>( &'a mut self, request: AppendRequest, placement: ShardPlacement, admission: ColdWriteAdmission, ) -> GroupAppendFuture<'a>
fn append_batch<'a>( &'a mut self, request: AppendBatchRequest, placement: ShardPlacement, admission: ColdWriteAdmission, ) -> GroupAppendBatchFuture<'a>
fn snapshot<'a>( &'a mut self, placement: ShardPlacement, ) -> GroupSnapshotFuture<'a>
fn install_snapshot<'a>( &'a mut self, snapshot: GroupSnapshot, ) -> GroupInstallSnapshotFuture<'a>
Provided Methods§
fn accepts_local_writes(&self) -> bool
fn create_stream_external<'a>( &'a mut self, request: CreateStreamExternalRequest, _placement: ShardPlacement, ) -> GroupCreateStreamFuture<'a>
fn get_stream_attrs<'a>( &'a mut self, request: GetStreamAttrsRequest, _placement: ShardPlacement, ) -> GroupGetStreamAttrsFuture<'a>
fn read_stream_parts<'a>( &'a mut self, request: ReadStreamRequest, placement: ShardPlacement, ) -> GroupReadStreamPartsFuture<'a>
fn require_local_live_read_owner<'a>( &'a mut self, _placement: ShardPlacement, ) -> GroupRequireLiveReadOwnerFuture<'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 update_stream_attrs<'a>( &'a mut self, request: UpdateStreamAttrsRequest, _placement: ShardPlacement, ) -> GroupUpdateStreamAttrsFuture<'a>
Sourcefn 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.
Sourcefn 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_external<'a>( &'a mut self, request: AppendExternalRequest, _placement: ShardPlacement, ) -> GroupAppendFuture<'a>
fn append_batch_many<'a>( &'a mut self, requests: Vec<AppendBatchRequest>, placement: ShardPlacement, admission: ColdWriteAdmission, ) -> GroupWriteBatchFuture<'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 shutdown<'a>(&'a mut self) -> GroupShutdownFuture<'a>
fn write_batch<'a>( &'a mut self, commands: Vec<GroupWriteCommand>, placement: ShardPlacement, ) -> GroupWriteBatchFuture<'a>
Sourcefn 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.
Dyn Compatibility§
This trait is dyn compatible.
In older versions of Rust, dyn compatibility was called "object safety".