pub trait GroupEngine: Send + 'static {
Show 38 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 bucket_usage<'a>(
&'a mut self,
placement: ShardPlacement,
) -> GroupBucketUsageFuture<'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 advance_retention<'a>(
&'a mut self,
request: AdvanceRetentionRequest,
_placement: ShardPlacement,
) -> GroupAdvanceRetentionFuture<'a> { ... }
fn import_group_state<'a>(
&'a mut self,
_request: ImportGroupStateRequest,
placement: ShardPlacement,
) -> GroupImportGroupStateFuture<'a> { ... }
fn set_bucket_quota<'a>(
&'a mut self,
request: SetBucketQuotaRequest,
_placement: ShardPlacement,
) -> GroupSetBucketQuotaFuture<'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 purge_bucket<'a>(
&'a mut self,
_bucket_id: String,
_placement: ShardPlacement,
) -> GroupPurgeBucketFuture<'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 append_transaction<'a>(
&'a mut self,
_request: AppendTransactionRequest,
_placement: ShardPlacement,
_admission: ColdWriteAdmission,
) -> GroupAppendTransactionFuture<'a> { ... }
fn flush_cold<'a>(
&'a mut self,
request: FlushColdRequest,
_placement: ShardPlacement,
) -> GroupFlushColdFuture<'a> { ... }
fn compact_cold<'a>(
&'a mut self,
request: CompactColdRequest,
_placement: ShardPlacement,
) -> GroupCompactColdFuture<'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>
Sourcefn bucket_usage<'a>(
&'a mut self,
placement: ShardPlacement,
) -> GroupBucketUsageFuture<'a>
fn bucket_usage<'a>( &'a mut self, placement: ShardPlacement, ) -> GroupBucketUsageFuture<'a>
Per-bucket committed usage held by this group’s replicated state.
Served from local replica state, leader or follower: usage export tolerates replication lag, and requiring leadership would make a node-local aggregate fail whenever any group is led elsewhere. Deliberately a required method — an engine that silently reported nothing would underbill.
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 advance_retention<'a>( &'a mut self, request: AdvanceRetentionRequest, _placement: ShardPlacement, ) -> GroupAdvanceRetentionFuture<'a>
Sourcefn import_group_state<'a>(
&'a mut self,
_request: ImportGroupStateRequest,
placement: ShardPlacement,
) -> GroupImportGroupStateFuture<'a>
fn import_group_state<'a>( &'a mut self, _request: ImportGroupStateRequest, placement: ShardPlacement, ) -> GroupImportGroupStateFuture<'a>
Restore path: replaces an empty group’s state with a backup snapshot
as one replicated write. See StreamCommand::ImportSnapshot.
fn set_bucket_quota<'a>( &'a mut self, request: SetBucketQuotaRequest, _placement: ShardPlacement, ) -> GroupSetBucketQuotaFuture<'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 purge_bucket<'a>(
&'a mut self,
_bucket_id: String,
_placement: ShardPlacement,
) -> GroupPurgeBucketFuture<'a>
fn purge_bucket<'a>( &'a mut self, _bucket_id: String, _placement: ShardPlacement, ) -> GroupPurgeBucketFuture<'a>
Replicated tenant offboarding: removes every stream in the bucket, the bucket, and its quota in this group. Monotonic aggregate usage remains available to asynchronous accounting readers. 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 append_transaction<'a>( &'a mut self, _request: AppendTransactionRequest, _placement: ShardPlacement, _admission: ColdWriteAdmission, ) -> GroupAppendTransactionFuture<'a>
fn flush_cold<'a>( &'a mut self, request: FlushColdRequest, _placement: ShardPlacement, ) -> GroupFlushColdFuture<'a>
fn compact_cold<'a>( &'a mut self, request: CompactColdRequest, _placement: ShardPlacement, ) -> GroupCompactColdFuture<'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".