use {
crate::{stream::BlockStream, wrapper::RESERVED_FILTER_NAME},
solana_commitment_config::CommitmentLevel,
tonic::async_trait,
yellowstone_grpc_client::{GeyserGrpcClient, GeyserGrpcClientError, GeyserStream},
yellowstone_grpc_proto::geyser::{
CommitmentLevel as ProtoCommitmentLevel, SubscribeRequest, SubscribeRequestFilterSlots,
SubscribeUpdate,
},
};
pub type GeyserBlockStream = BlockStream<GeyserStream, SubscribeUpdate>;
#[async_trait]
pub trait GeyserGrpcExt {
async fn subscribe_block(
&mut self,
subscribe_request: SubscribeRequest,
) -> Result<GeyserBlockStream, GeyserGrpcClientError>;
}
pub const DEFAULT_SUBSCRIBE_BLOCK_CHANNEL_CAPACITY: usize = 1_000_000;
#[derive(Debug, thiserror::Error)]
pub enum BlockMachineError {
#[error(transparent)]
GrpcError(#[from] tonic::Status),
}
#[async_trait]
impl GeyserGrpcExt for GeyserGrpcClient {
async fn subscribe_block(
&mut self,
mut subscribe_request: SubscribeRequest,
) -> Result<GeyserBlockStream, GeyserGrpcClientError> {
let proto_commitment_level =
ProtoCommitmentLevel::try_from(subscribe_request.commitment.unwrap_or(0))
.expect("Invalid commitment level in subscribe request");
assert!(
subscribe_request.blocks.is_empty(),
"custom `blocks` filter is not compatible with block machine"
);
assert!(
subscribe_request.slots.is_empty(),
"custom `slots` filter is not compatible with block machine"
);
let commitment_level = match proto_commitment_level {
ProtoCommitmentLevel::Processed => CommitmentLevel::Processed,
ProtoCommitmentLevel::Confirmed => CommitmentLevel::Confirmed,
ProtoCommitmentLevel::Finalized => CommitmentLevel::Finalized,
};
subscribe_request.slots.insert(
RESERVED_FILTER_NAME.to_owned(),
SubscribeRequestFilterSlots {
interslot_updates: Some(true),
..Default::default()
},
);
subscribe_request
.blocks_meta
.insert(RESERVED_FILTER_NAME.to_owned(), Default::default());
subscribe_request
.entry
.insert(RESERVED_FILTER_NAME.to_owned(), Default::default());
subscribe_request.commitment = Some(0);
let (_sink, source) = self.subscribe_with_request(Some(subscribe_request)).await?;
Ok(BlockStream::new(source, commitment_level))
}
}