macro_rules! runtime_operations {
($generate:ident) => {
$generate! {
op CreateStream {
fields { request: CreateStreamRequest }
reply { response_tx: CreateStreamResponse }
guard { raft_uncommitted }
handle { call create_stream(engine, metrics, request, placement, cold_admission) }
client {
pub stream fn create_stream,
admit: u64::try_from(request.initial_payload.len())
.expect("payload len fits u64")
}
}
op CreateExternal {
fields { request: CreateStreamExternalRequest }
reply { response_tx: CreateStreamResponse }
guard { none }
handle { call create_stream_external(engine, metrics, request, placement) }
client { pub stream fn create_stream_external }
}
op HeadStream {
fields { request: HeadStreamRequest }
reply { response_tx: HeadStreamResponse }
guard { none }
handle { call head_stream(engine, metrics, request, placement) }
client { pub stream fn head_stream }
}
op GetStreamAttrs {
fields { request: GetStreamAttrsRequest }
reply { response_tx: GetStreamAttrsResponse }
guard { none }
handle { call get_stream_attrs(engine, metrics, request, placement) }
client { pub stream fn get_stream_attrs }
}
op ReadStream {
fields { request: ReadStreamRequest }
reply { response_tx: ReadStreamResponse }
guard { none }
handle {
tail read_stream(
engine,
metrics,
read_materialization,
request,
placement,
response_tx
)
}
client { pub stream fn read_stream }
}
op PublishSnapshot {
fields { request: PublishSnapshotRequest }
reply { response_tx: PublishSnapshotResponse }
guard { none }
handle {
call publish_snapshot(
engine,
metrics,
read_materialization,
read_watchers,
request,
placement
)
}
client { pub stream fn publish_snapshot }
}
op AdvanceRetention {
fields { request: AdvanceRetentionRequest }
reply { response_tx: AdvanceRetentionResponse }
guard { none }
handle {
call advance_retention(
engine,
metrics,
read_materialization,
read_watchers,
request,
placement
)
}
client { pub stream fn advance_retention }
}
op ReadSnapshot {
fields { request: ReadSnapshotRequest }
reply { response_tx: ReadSnapshotResponse }
guard { none }
handle { call read_snapshot(engine, metrics, request, placement) }
client { pub stream fn read_snapshot }
}
op DeleteSnapshot {
fields { request: DeleteSnapshotRequest }
reply { response_tx: () }
guard { none }
handle { call delete_snapshot(engine, metrics, request, placement) }
client { pub stream fn delete_snapshot }
}
op BootstrapStream {
fields { request: BootstrapStreamRequest }
reply { response_tx: BootstrapStreamResponse }
guard { none }
handle { call bootstrap_stream(engine, metrics, request, placement) }
client { pub stream fn bootstrap_stream }
}
op WaitRead {
fields { request: ReadStreamRequest, waiter_id: u64 }
reply { response_tx: ReadStreamResponse }
guard { none }
handle { actor handle_wait_read(request, waiter_id, response_tx) }
client { none }
}
op CancelWaitRead {
fields { stream_id: BucketStreamId, waiter_id: u64 }
reply { none }
guard { none }
handle {
sync cancel_read_watcher(read_watchers, metrics, core_id, stream_id, waiter_id)
}
client { none }
}
op RequireLiveReadOwner {
fields {}
reply { response_tx: () }
guard { none }
handle { call require_live_read_owner(engine, placement) }
client { none }
}
op CloseStream {
fields { request: CloseStreamRequest }
reply { response_tx: CloseStreamResponse }
guard { none }
handle {
call close_stream(
engine,
metrics,
read_materialization,
read_watchers,
request,
placement
)
}
client { pub stream fn close_stream }
}
op UpdateStreamAttrs {
fields { request: UpdateStreamAttrsRequest }
reply { response_tx: UpdateStreamAttrsResponse }
guard { none }
handle { call update_stream_attrs(engine, metrics, request, placement) }
client { pub stream fn update_stream_attrs }
}
op DeleteStream {
fields { request: DeleteStreamRequest }
reply { response_tx: DeleteStreamResponse }
guard { none }
handle {
call delete_stream(
engine,
metrics,
read_materialization,
read_watchers,
request,
placement
)
}
client { pub stream fn delete_stream }
}
op FlushCold {
fields { request: FlushColdRequest }
reply { response_tx: FlushColdResponse }
guard { none }
handle {
call flush_cold(
engine,
metrics,
read_materialization,
read_watchers,
request,
placement
)
}
client { pub stream fn flush_cold }
}
op CompactCold {
fields { request: CompactColdRequest }
reply { response_tx: CompactColdResponse }
guard { none }
handle {
call compact_cold(
engine,
metrics,
read_materialization,
read_watchers,
request,
placement
)
}
client { pub stream fn compact_cold }
}
op PlanColdFlush {
fields { request: PlanColdFlushRequest }
reply { response_tx: Option<ColdFlushCandidate> }
guard { none }
handle { call plan_cold_flush(engine, metrics, request, placement) }
client { pub stream fn plan_cold_flush }
}
op PlanNextColdFlushBatch {
fields { request: PlanGroupColdFlushRequest, max_candidates: usize }
reply { response_tx: Vec<ColdFlushCandidate> }
guard { none }
handle {
call plan_next_cold_flush_batch(
engine,
metrics,
request,
placement,
max_candidates
)
}
client { pub group fn plan_next_cold_flush_batch }
}
op PlanColdGc {
fields { max: usize }
reply { response_tx: Vec<ColdGcEntry> }
guard { none }
handle { call plan_cold_gc(engine, max, placement) }
client { group fn plan_cold_gc }
}
op BucketUsage {
fields {}
reply { response_tx: Vec<ursula_stream::BucketUsageSnapshot> }
guard { none }
handle { call bucket_usage(engine, metrics, placement) }
client { pub group fn bucket_usage }
}
op SetBucketQuota {
fields { request: SetBucketQuotaRequest }
reply { response_tx: SetBucketQuotaResponse }
guard { none }
handle { call set_bucket_quota(engine, metrics, request, placement) }
client { pub group fn set_bucket_quota }
}
op AckColdGc {
fields { up_to_seq: u64 }
reply { response_tx: AckColdGcResponse }
guard { none }
handle { call ack_cold_gc(engine, up_to_seq, placement) }
client { group fn ack_cold_gc }
}
op PurgeBucket {
fields { bucket_id: String }
reply { response_tx: PurgeBucketResponse }
guard { none }
handle { call purge_bucket(engine, bucket_id, placement) }
client { pub group fn purge_bucket }
}
op Append {
fields { request: AppendRequest }
reply { response_tx: AppendResponse }
guard { raft_uncommitted }
handle { actor handle_append(request, response_tx, raft_uncommitted) }
client {
pub stream fn append,
non_empty: payload,
admit: request.payload_len()
}
}
op AppendTransaction {
fields { request: AppendTransactionRequest }
reply { response_tx: AppendTransactionResponse }
guard { raft_uncommitted }
handle { actor handle_append_transaction(request, response_tx, raft_uncommitted) }
client { none }
}
op AppendExternal {
fields { request: AppendExternalRequest }
reply { response_tx: AppendResponse }
guard { none }
handle {
call apply_append_external(
engine,
metrics,
read_materialization,
read_watchers,
request,
placement
)
}
client { pub stream fn append_external }
}
op AppendBatch {
fields { request: AppendBatchRequest }
reply { response_tx: AppendBatchResponse }
guard { raft_uncommitted }
handle { actor handle_append_batch(request, response_tx, raft_uncommitted, pending) }
client {
pub stream fn append_batch,
non_empty: payloads,
admit: append_batch_payload_bytes(&request)
}
}
op SnapshotGroup {
fields {}
reply { response_tx: GroupSnapshot }
guard { none }
handle { call snapshot_group(engine, metrics, placement) }
client { pub group fn snapshot_group }
}
op ImportGroupState {
fields { request: ImportGroupStateRequest }
reply { response_tx: ImportGroupStateResponse }
guard { none }
handle { call import_group_state(engine, metrics, request, placement) }
client { pub group fn import_group_state }
}
op InstallGroupSnapshot {
fields { snapshot: GroupSnapshot }
reply { response_tx: () }
guard { none }
handle { call install_group_snapshot(engine, metrics, snapshot) }
client { none }
}
op ShutdownEngine {
fields {}
reply { response_tx: () }
guard { none }
handle { actor handle_shutdown_engine(response_tx) }
client { none }
}
}
};
}
pub(crate) use runtime_operations;