rebuild code

This commit is contained in:
wood
2025-07-06 22:06:44 +08:00
parent 9e202a1e5e
commit 3ac39508d9
34 changed files with 635 additions and 552 deletions
+260
View File
@@ -0,0 +1,260 @@
// This file is @generated by prost-build.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct GenerateAuthChallengeRequest {
/// / Role the client is attempting to generate tokens for.
#[prost(enumeration = "Role", tag = "1")]
pub role: i32,
/// / Client's 32 byte pubkey.
#[prost(bytes = "vec", tag = "2")]
pub pubkey: ::prost::alloc::vec::Vec<u8>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct GenerateAuthChallengeResponse {
#[prost(string, tag = "1")]
pub challenge: ::prost::alloc::string::String,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct GenerateAuthTokensRequest {
/// / The pre-signed challenge.
#[prost(string, tag = "1")]
pub challenge: ::prost::alloc::string::String,
/// / The signing keypair's corresponding 32 byte pubkey.
#[prost(bytes = "vec", tag = "2")]
pub client_pubkey: ::prost::alloc::vec::Vec<u8>,
/// / The 64 byte signature of the challenge signed by the client's private key. The private key must correspond to
/// the pubkey passed in the \[GenerateAuthChallenge\] method. The client is expected to sign the challenge token
/// prepended with their pubkey. For example sign(pubkey, challenge).
#[prost(bytes = "vec", tag = "3")]
pub signed_challenge: ::prost::alloc::vec::Vec<u8>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Token {
/// / The token.
#[prost(string, tag = "1")]
pub value: ::prost::alloc::string::String,
/// / When the token will expire.
#[prost(message, optional, tag = "2")]
pub expires_at_utc: ::core::option::Option<::prost_types::Timestamp>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct GenerateAuthTokensResponse {
/// / The token granting access to resources.
#[prost(message, optional, tag = "1")]
pub access_token: ::core::option::Option<Token>,
/// / The token used to refresh the access_token. This has a longer TTL than the access_token.
#[prost(message, optional, tag = "2")]
pub refresh_token: ::core::option::Option<Token>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct RefreshAccessTokenRequest {
/// / Non-expired refresh token obtained from the \[GenerateAuthTokens\] method.
#[prost(string, tag = "1")]
pub refresh_token: ::prost::alloc::string::String,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct RefreshAccessTokenResponse {
/// / Fresh access_token.
#[prost(message, optional, tag = "1")]
pub access_token: ::core::option::Option<Token>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
#[repr(i32)]
pub enum Role {
Relayer = 0,
Searcher = 1,
Validator = 2,
ShredstreamSubscriber = 3,
}
impl Role {
/// String value of the enum field names used in the ProtoBuf definition.
///
/// The values are not transformed in any way and thus are considered stable
/// (if the ProtoBuf definition does not change) and safe for programmatic use.
pub fn as_str_name(&self) -> &'static str {
match self {
Self::Relayer => "RELAYER",
Self::Searcher => "SEARCHER",
Self::Validator => "VALIDATOR",
Self::ShredstreamSubscriber => "SHREDSTREAM_SUBSCRIBER",
}
}
/// Creates an enum from field names used in the ProtoBuf definition.
pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
match value {
"RELAYER" => Some(Self::Relayer),
"SEARCHER" => Some(Self::Searcher),
"VALIDATOR" => Some(Self::Validator),
"SHREDSTREAM_SUBSCRIBER" => Some(Self::ShredstreamSubscriber),
_ => None,
}
}
}
/// Generated client implementations.
pub mod auth_service_client {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
use tonic::codegen::http::Uri;
/// / This service is responsible for issuing auth tokens to clients for API access.
#[derive(Debug, Clone)]
pub struct AuthServiceClient<T> {
inner: tonic::client::Grpc<T>,
}
impl AuthServiceClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> AuthServiceClient<T>
where
T: tonic::client::GrpcService<tonic::body::BoxBody>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
{
pub fn new(inner: T) -> Self {
let inner = tonic::client::Grpc::new(inner);
Self { inner }
}
pub fn with_origin(inner: T, origin: Uri) -> Self {
let inner = tonic::client::Grpc::with_origin(inner, origin);
Self { inner }
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> AuthServiceClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
Response = http::Response<
<T as tonic::client::GrpcService<tonic::body::BoxBody>>::ResponseBody,
>,
>,
<T as tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
>>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
{
AuthServiceClient::new(InterceptedService::new(inner, interceptor))
}
/// Compress requests with the given encoding.
///
/// This requires the server to support it otherwise it might respond with an
/// error.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.send_compressed(encoding);
self
}
/// Enable decompressing responses.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.accept_compressed(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_decoding_message_size(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_encoding_message_size(limit);
self
}
/// / Returns a challenge, client is expected to sign this challenge with an appropriate keypair in order to obtain access tokens.
pub async fn generate_auth_challenge(
&mut self,
request: impl tonic::IntoRequest<super::GenerateAuthChallengeRequest>,
) -> std::result::Result<
tonic::Response<super::GenerateAuthChallengeResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/auth.AuthService/GenerateAuthChallenge",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("auth.AuthService", "GenerateAuthChallenge"));
self.inner.unary(req, path, codec).await
}
/// / Provides the client with the initial pair of auth tokens for API access.
pub async fn generate_auth_tokens(
&mut self,
request: impl tonic::IntoRequest<super::GenerateAuthTokensRequest>,
) -> std::result::Result<
tonic::Response<super::GenerateAuthTokensResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/auth.AuthService/GenerateAuthTokens",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("auth.AuthService", "GenerateAuthTokens"));
self.inner.unary(req, path, codec).await
}
/// / Call this method with a non-expired refresh token to obtain a new access token.
pub async fn refresh_access_token(
&mut self,
request: impl tonic::IntoRequest<super::RefreshAccessTokenRequest>,
) -> std::result::Result<
tonic::Response<super::RefreshAccessTokenResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/auth.AuthService/RefreshAccessToken",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("auth.AuthService", "RefreshAccessToken"));
self.inner.unary(req, path, codec).await
}
}
}
+19
View File
@@ -0,0 +1,19 @@
// This file is @generated by prost-build.
/// Condensed block helpful for getting data around efficiently internal to our system.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct CondensedBlock {
#[prost(message, optional, tag = "1")]
pub header: ::core::option::Option<super::shared::Header>,
#[prost(string, tag = "2")]
pub previous_blockhash: ::prost::alloc::string::String,
#[prost(string, tag = "3")]
pub blockhash: ::prost::alloc::string::String,
#[prost(uint64, tag = "4")]
pub parent_slot: u64,
#[prost(bytes = "vec", repeated, tag = "5")]
pub versioned_transactions: ::prost::alloc::vec::Vec<::prost::alloc::vec::Vec<u8>>,
#[prost(uint64, tag = "6")]
pub slot: u64,
#[prost(string, tag = "7")]
pub commitment: ::prost::alloc::string::String,
}
+462
View File
@@ -0,0 +1,462 @@
// This file is @generated by prost-build.
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct SubscribePacketsRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct SubscribePacketsResponse {
#[prost(message, optional, tag = "1")]
pub header: ::core::option::Option<super::shared::Header>,
#[prost(message, optional, tag = "2")]
pub batch: ::core::option::Option<super::packet::PacketBatch>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct SubscribeBundlesRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct SubscribeBundlesResponse {
#[prost(message, repeated, tag = "1")]
pub bundles: ::prost::alloc::vec::Vec<super::bundle::BundleUuid>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct BlockBuilderFeeInfoRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct BlockBuilderFeeInfoResponse {
#[prost(string, tag = "1")]
pub pubkey: ::prost::alloc::string::String,
/// commission (0-100)
#[prost(uint64, tag = "2")]
pub commission: u64,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct AccountsOfInterest {
/// use * for all accounts
#[prost(string, repeated, tag = "1")]
pub accounts: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct AccountsOfInterestRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct AccountsOfInterestUpdate {
#[prost(string, repeated, tag = "1")]
pub accounts: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct ProgramsOfInterestRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct ProgramsOfInterestUpdate {
#[prost(string, repeated, tag = "1")]
pub programs: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
}
/// A series of packets with an expiration attached to them.
/// The header contains a timestamp for when this packet was generated.
/// The expiry is how long the packet batches have before they expire and are forwarded to the validator.
/// This provides a more censorship resistant method to MEV than block engines receiving packets directly.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct ExpiringPacketBatch {
#[prost(message, optional, tag = "1")]
pub header: ::core::option::Option<super::shared::Header>,
#[prost(message, optional, tag = "2")]
pub batch: ::core::option::Option<super::packet::PacketBatch>,
#[prost(uint32, tag = "3")]
pub expiry_ms: u32,
}
/// Packets and heartbeats are sent over the same stream.
/// ExpiringPacketBatches have an expiration attached to them so the block engine can track
/// how long it has until the relayer forwards the packets to the validator.
/// Heartbeats contain a timestamp from the system and is used as a simple and naive time-sync mechanism
/// so the block engine has some idea on how far their clocks are apart.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct PacketBatchUpdate {
#[prost(oneof = "packet_batch_update::Msg", tags = "1, 2")]
pub msg: ::core::option::Option<packet_batch_update::Msg>,
}
/// Nested message and enum types in `PacketBatchUpdate`.
pub mod packet_batch_update {
#[derive(Clone, PartialEq, ::prost::Oneof)]
pub enum Msg {
#[prost(message, tag = "1")]
Batches(super::ExpiringPacketBatch),
#[prost(message, tag = "2")]
Heartbeat(super::super::shared::Heartbeat),
}
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct StartExpiringPacketStreamResponse {
#[prost(message, optional, tag = "1")]
pub heartbeat: ::core::option::Option<super::shared::Heartbeat>,
}
/// Generated client implementations.
pub mod block_engine_validator_client {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
use tonic::codegen::http::Uri;
/// / Validators can connect to Block Engines to receive packets and bundles.
#[derive(Debug, Clone)]
pub struct BlockEngineValidatorClient<T> {
inner: tonic::client::Grpc<T>,
}
impl BlockEngineValidatorClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> BlockEngineValidatorClient<T>
where
T: tonic::client::GrpcService<tonic::body::BoxBody>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
{
pub fn new(inner: T) -> Self {
let inner = tonic::client::Grpc::new(inner);
Self { inner }
}
pub fn with_origin(inner: T, origin: Uri) -> Self {
let inner = tonic::client::Grpc::with_origin(inner, origin);
Self { inner }
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> BlockEngineValidatorClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
Response = http::Response<
<T as tonic::client::GrpcService<tonic::body::BoxBody>>::ResponseBody,
>,
>,
<T as tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
>>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
{
BlockEngineValidatorClient::new(InterceptedService::new(inner, interceptor))
}
/// Compress requests with the given encoding.
///
/// This requires the server to support it otherwise it might respond with an
/// error.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.send_compressed(encoding);
self
}
/// Enable decompressing responses.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.accept_compressed(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_decoding_message_size(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_encoding_message_size(limit);
self
}
/// / Validators can subscribe to the block engine to receive a stream of packets
pub async fn subscribe_packets(
&mut self,
request: impl tonic::IntoRequest<super::SubscribePacketsRequest>,
) -> std::result::Result<
tonic::Response<tonic::codec::Streaming<super::SubscribePacketsResponse>>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/block_engine.BlockEngineValidator/SubscribePackets",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new(
"block_engine.BlockEngineValidator",
"SubscribePackets",
),
);
self.inner.server_streaming(req, path, codec).await
}
/// / Validators can subscribe to the block engine to receive a stream of simulated and profitable bundles
pub async fn subscribe_bundles(
&mut self,
request: impl tonic::IntoRequest<super::SubscribeBundlesRequest>,
) -> std::result::Result<
tonic::Response<tonic::codec::Streaming<super::SubscribeBundlesResponse>>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/block_engine.BlockEngineValidator/SubscribeBundles",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new(
"block_engine.BlockEngineValidator",
"SubscribeBundles",
),
);
self.inner.server_streaming(req, path, codec).await
}
/// Block builders can optionally collect fees. This returns fee information if a block builder wants to
/// collect one.
pub async fn get_block_builder_fee_info(
&mut self,
request: impl tonic::IntoRequest<super::BlockBuilderFeeInfoRequest>,
) -> std::result::Result<
tonic::Response<super::BlockBuilderFeeInfoResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/block_engine.BlockEngineValidator/GetBlockBuilderFeeInfo",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new(
"block_engine.BlockEngineValidator",
"GetBlockBuilderFeeInfo",
),
);
self.inner.unary(req, path, codec).await
}
}
}
/// Generated client implementations.
pub mod block_engine_relayer_client {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
use tonic::codegen::http::Uri;
/// / Relayers can forward packets to Block Engines.
/// / Block Engines provide an AccountsOfInterest field to only send transactions that are of interest.
#[derive(Debug, Clone)]
pub struct BlockEngineRelayerClient<T> {
inner: tonic::client::Grpc<T>,
}
impl BlockEngineRelayerClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> BlockEngineRelayerClient<T>
where
T: tonic::client::GrpcService<tonic::body::BoxBody>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
{
pub fn new(inner: T) -> Self {
let inner = tonic::client::Grpc::new(inner);
Self { inner }
}
pub fn with_origin(inner: T, origin: Uri) -> Self {
let inner = tonic::client::Grpc::with_origin(inner, origin);
Self { inner }
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> BlockEngineRelayerClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
Response = http::Response<
<T as tonic::client::GrpcService<tonic::body::BoxBody>>::ResponseBody,
>,
>,
<T as tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
>>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
{
BlockEngineRelayerClient::new(InterceptedService::new(inner, interceptor))
}
/// Compress requests with the given encoding.
///
/// This requires the server to support it otherwise it might respond with an
/// error.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.send_compressed(encoding);
self
}
/// Enable decompressing responses.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.accept_compressed(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_decoding_message_size(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_encoding_message_size(limit);
self
}
/// / The block engine feeds accounts of interest (AOI) updates to the relayer periodically.
/// / For all transactions the relayer receives, it forwards transactions to the block engine which write-lock
/// / any of the accounts in the AOI.
pub async fn subscribe_accounts_of_interest(
&mut self,
request: impl tonic::IntoRequest<super::AccountsOfInterestRequest>,
) -> std::result::Result<
tonic::Response<tonic::codec::Streaming<super::AccountsOfInterestUpdate>>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/block_engine.BlockEngineRelayer/SubscribeAccountsOfInterest",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new(
"block_engine.BlockEngineRelayer",
"SubscribeAccountsOfInterest",
),
);
self.inner.server_streaming(req, path, codec).await
}
pub async fn subscribe_programs_of_interest(
&mut self,
request: impl tonic::IntoRequest<super::ProgramsOfInterestRequest>,
) -> std::result::Result<
tonic::Response<tonic::codec::Streaming<super::ProgramsOfInterestUpdate>>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/block_engine.BlockEngineRelayer/SubscribeProgramsOfInterest",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new(
"block_engine.BlockEngineRelayer",
"SubscribeProgramsOfInterest",
),
);
self.inner.server_streaming(req, path, codec).await
}
/// Validators can subscribe to packets from the relayer and receive a multiplexed signal that contains a mixture
/// of packets and heartbeats.
/// NOTE: This is a bi-directional stream due to a bug with how Envoy handles half closed client-side streams.
/// The issue is being tracked here: https://github.com/envoyproxy/envoy/issues/22748. In the meantime, the
/// server will stream heartbeats to clients at some reasonable cadence.
pub async fn start_expiring_packet_stream(
&mut self,
request: impl tonic::IntoStreamingRequest<Message = super::PacketBatchUpdate>,
) -> std::result::Result<
tonic::Response<
tonic::codec::Streaming<super::StartExpiringPacketStreamResponse>,
>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/block_engine.BlockEngineRelayer/StartExpiringPacketStream",
);
let mut req = request.into_streaming_request();
req.extensions_mut()
.insert(
GrpcMethod::new(
"block_engine.BlockEngineRelayer",
"StartExpiringPacketStream",
),
);
self.inner.streaming(req, path, codec).await
}
}
}
+171
View File
@@ -0,0 +1,171 @@
// This file is @generated by prost-build.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Bundle {
#[prost(message, optional, tag = "2")]
pub header: ::core::option::Option<super::shared::Header>,
#[prost(message, repeated, tag = "3")]
pub packets: ::prost::alloc::vec::Vec<super::packet::Packet>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct BundleUuid {
#[prost(message, optional, tag = "1")]
pub bundle: ::core::option::Option<Bundle>,
#[prost(string, tag = "2")]
pub uuid: ::prost::alloc::string::String,
}
/// Indicates the bundle was accepted and forwarded to a validator.
/// NOTE: A single bundle may have multiple events emitted if forwarded to many validators.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Accepted {
/// Slot at which bundle was forwarded.
#[prost(uint64, tag = "1")]
pub slot: u64,
/// Validator identity bundle was forwarded to.
#[prost(string, tag = "2")]
pub validator_identity: ::prost::alloc::string::String,
}
/// Indicates the bundle was dropped and therefore not forwarded to any validator.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Rejected {
#[prost(oneof = "rejected::Reason", tags = "1, 2, 3, 4, 5")]
pub reason: ::core::option::Option<rejected::Reason>,
}
/// Nested message and enum types in `Rejected`.
pub mod rejected {
#[derive(Clone, PartialEq, ::prost::Oneof)]
pub enum Reason {
#[prost(message, tag = "1")]
StateAuctionBidRejected(super::StateAuctionBidRejected),
#[prost(message, tag = "2")]
WinningBatchBidRejected(super::WinningBatchBidRejected),
#[prost(message, tag = "3")]
SimulationFailure(super::SimulationFailure),
#[prost(message, tag = "4")]
InternalError(super::InternalError),
#[prost(message, tag = "5")]
DroppedBundle(super::DroppedBundle),
}
}
/// Indicates the bundle's bid was high enough to win its state auction.
/// However, not high enough relative to other state auction winners and therefore excluded from being forwarded.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct WinningBatchBidRejected {
/// Auction's unique identifier.
#[prost(string, tag = "1")]
pub auction_id: ::prost::alloc::string::String,
/// Bundle's simulated bid.
#[prost(uint64, tag = "2")]
pub simulated_bid_lamports: u64,
#[prost(string, optional, tag = "3")]
pub msg: ::core::option::Option<::prost::alloc::string::String>,
}
/// Indicates the bundle's bid was __not__ high enough to be included in its state auction's set of winners.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct StateAuctionBidRejected {
/// Auction's unique identifier.
#[prost(string, tag = "1")]
pub auction_id: ::prost::alloc::string::String,
/// Bundle's simulated bid.
#[prost(uint64, tag = "2")]
pub simulated_bid_lamports: u64,
#[prost(string, optional, tag = "3")]
pub msg: ::core::option::Option<::prost::alloc::string::String>,
}
/// Bundle dropped due to simulation failure.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct SimulationFailure {
/// Signature of the offending transaction.
#[prost(string, tag = "1")]
pub tx_signature: ::prost::alloc::string::String,
#[prost(string, optional, tag = "2")]
pub msg: ::core::option::Option<::prost::alloc::string::String>,
}
/// Bundle dropped due to an internal error.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct InternalError {
#[prost(string, tag = "1")]
pub msg: ::prost::alloc::string::String,
}
/// Bundle dropped (e.g. because no leader upcoming)
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct DroppedBundle {
#[prost(string, tag = "1")]
pub msg: ::prost::alloc::string::String,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct Finalized {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Processed {
#[prost(string, tag = "1")]
pub validator_identity: ::prost::alloc::string::String,
#[prost(uint64, tag = "2")]
pub slot: u64,
/// / Index within the block.
#[prost(uint64, tag = "3")]
pub bundle_index: u64,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct Dropped {
#[prost(enumeration = "DroppedReason", tag = "1")]
pub reason: i32,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct BundleResult {
/// Bundle's Uuid.
#[prost(string, tag = "1")]
pub bundle_id: ::prost::alloc::string::String,
#[prost(oneof = "bundle_result::Result", tags = "2, 3, 4, 5, 6")]
pub result: ::core::option::Option<bundle_result::Result>,
}
/// Nested message and enum types in `BundleResult`.
pub mod bundle_result {
#[derive(Clone, PartialEq, ::prost::Oneof)]
pub enum Result {
/// Indicated accepted by the block-engine and forwarded to a jito-solana validator.
#[prost(message, tag = "2")]
Accepted(super::Accepted),
/// Rejected by the block-engine.
#[prost(message, tag = "3")]
Rejected(super::Rejected),
/// Reached finalized commitment level.
#[prost(message, tag = "4")]
Finalized(super::Finalized),
/// Reached a processed commitment level.
#[prost(message, tag = "5")]
Processed(super::Processed),
/// Was accepted and forwarded by the block-engine but never landed on-chain.
#[prost(message, tag = "6")]
Dropped(super::Dropped),
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
#[repr(i32)]
pub enum DroppedReason {
BlockhashExpired = 0,
/// One or more transactions in the bundle landed on-chain, invalidating the bundle.
PartiallyProcessed = 1,
/// This indicates bundle was processed but not finalized. This could occur during forks.
NotFinalized = 2,
}
impl DroppedReason {
/// String value of the enum field names used in the ProtoBuf definition.
///
/// The values are not transformed in any way and thus are considered stable
/// (if the ProtoBuf definition does not change) and safe for programmatic use.
pub fn as_str_name(&self) -> &'static str {
match self {
Self::BlockhashExpired => "BlockhashExpired",
Self::PartiallyProcessed => "PartiallyProcessed",
Self::NotFinalized => "NotFinalized",
}
}
/// Creates an enum from field names used in the ProtoBuf definition.
pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
match value {
"BlockhashExpired" => Some(Self::BlockhashExpired),
"PartiallyProcessed" => Some(Self::PartiallyProcessed),
"NotFinalized" => Some(Self::NotFinalized),
_ => None,
}
}
}
+158
View File
@@ -0,0 +1,158 @@
use std::{
cmp::min,
net::{AddrParseError, IpAddr, Ipv4Addr, SocketAddr},
str::FromStr,
};
use bincode::serialize;
use solana_perf::packet::{Packet, PacketBatch, PACKET_DATA_SIZE};
use solana_sdk::{
packet::{Meta, PacketFlags},
transaction::VersionedTransaction,
};
use crate::protos::{
packet::{
Meta as ProtoMeta, Packet as ProtoPacket, PacketBatch as ProtoPacketBatch,
PacketFlags as ProtoPacketFlags,
},
shared::Socket,
};
/// Converts a Solana packet to a protobuf packet
/// NOTE: the packet.data() function will filter packets marked for discard
pub fn packet_to_proto_packet(p: &Packet) -> Option<ProtoPacket> {
Some(ProtoPacket {
data: p.data(..)?.to_vec(),
meta: Some(ProtoMeta {
size: p.meta().size as u64,
addr: p.meta().addr.to_string(),
port: p.meta().port as u32,
flags: Some(ProtoPacketFlags {
discard: p.meta().discard(),
forwarded: p.meta().forwarded(),
repair: p.meta().repair(),
simple_vote_tx: p.meta().is_simple_vote_tx(),
tracer_packet: p.meta().is_perf_track_packet(),
from_staked_node: p.meta().is_from_staked_node(),
}),
sender_stake: 0,
}),
})
}
pub fn packet_batches_to_proto_packets(
batches: &[PacketBatch],
) -> impl Iterator<Item = ProtoPacket> + '_ {
batches
.iter()
.flat_map(|b| b.iter().filter_map(packet_to_proto_packet))
}
/// converts from a protobuf packet to packet
pub fn proto_packet_to_packet(p: &ProtoPacket) -> Packet {
let mut data = [0u8; PACKET_DATA_SIZE];
let copy_len = min(data.len(), p.data.len());
data[..copy_len].copy_from_slice(&p.data[..copy_len]);
let mut packet = Packet::new(data, Meta::default());
if let Some(meta) = &p.meta {
packet.meta_mut().size = meta.size as usize;
packet.meta_mut().addr = meta
.addr
.parse()
.unwrap_or(IpAddr::V4(Ipv4Addr::UNSPECIFIED));
packet.meta_mut().port = meta.port as u16;
if let Some(flags) = &meta.flags {
if flags.simple_vote_tx {
packet.meta_mut().flags.insert(PacketFlags::SIMPLE_VOTE_TX);
}
if flags.forwarded {
packet.meta_mut().flags.insert(PacketFlags::FORWARDED);
}
if flags.tracer_packet {
packet.meta_mut().flags.insert(PacketFlags::PERF_TRACK_PACKET);
}
if flags.repair {
packet.meta_mut().flags.insert(PacketFlags::REPAIR);
}
if flags.discard {
packet.meta_mut().flags.insert(PacketFlags::DISCARD);
}
}
}
packet
}
pub fn proto_packet_batch_to_packets(
packet_batch: ProtoPacketBatch,
) -> impl Iterator<Item = Packet> {
packet_batch
.packets
.into_iter()
.map(|proto_packet| proto_packet_to_packet(&proto_packet))
}
/// Converts a protobuf packet to a VersionedTransaction
pub fn versioned_tx_from_packet(p: &ProtoPacket) -> Option<VersionedTransaction> {
let mut data = [0; PACKET_DATA_SIZE];
let copy_len = min(data.len(), p.data.len());
data[..copy_len].copy_from_slice(&p.data[..copy_len]);
let mut packet = Packet::new(data, Default::default());
if let Some(meta) = &p.meta {
packet.meta_mut().size = meta.size as usize;
}
packet.deserialize_slice(..).ok()
}
/// Coverts a VersionedTransaction to packet
pub fn packet_from_versioned_tx(tx: VersionedTransaction) -> Packet {
let tx_data = serialize(&tx).expect("serializes");
let mut data = [0; PACKET_DATA_SIZE];
let copy_len = min(tx_data.len(), data.len());
data[..copy_len].copy_from_slice(&tx_data[..copy_len]);
let mut packet = Packet::new(data, Default::default());
packet.meta_mut().size = copy_len;
packet
}
/// Converts a VersionedTransaction to a protobuf packet
pub fn proto_packet_from_versioned_tx(tx: &VersionedTransaction) -> ProtoPacket {
let data = serialize(tx).expect("serializes");
let size = data.len() as u64;
ProtoPacket {
data,
meta: Some(ProtoMeta {
size,
addr: "".to_string(),
port: 0,
flags: None,
sender_stake: 0,
}),
}
}
/// Converts a GRPC Socket to stdlib SocketAddr
impl TryFrom<&Socket> for SocketAddr {
type Error = AddrParseError;
fn try_from(value: &Socket) -> Result<Self, Self::Error> {
IpAddr::from_str(&value.ip).map(|ip| SocketAddr::new(ip, value.port as u16))
}
}
// #[cfg(test)]
// mod tests {
// use solana_perf::test_tx::test_tx;
// use solana_sdk::transaction::VersionedTransaction;
// use crate::convert::{proto_packet_from_versioned_tx, versioned_tx_from_packet};
// #[test]
// fn test_proto_to_packet() {
// let tx_before = VersionedTransaction::from(test_tx());
// let tx_after = versioned_tx_from_packet(&proto_packet_from_versioned_tx(&tx_before))
// .expect("tx_after");
// assert_eq!(tx_before, tx_after);
// }
// }
+14
View File
@@ -0,0 +1,14 @@
pub mod auth;
pub mod block;
pub mod block_engine;
pub mod bundle;
pub mod packet;
pub mod relayer;
pub mod searcher;
pub mod shared;
pub mod shredstream;
pub mod trace_shred;
pub mod convert;
pub mod nextblock_grpc;
pub mod searcher_client;
pub mod token_authenticator;
+462
View File
@@ -0,0 +1,462 @@
// This file is @generated by prost-build.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct PostSubmitRequest {
#[prost(message, optional, tag = "1")]
pub transaction: ::core::option::Option<TransactionMessage>,
#[prost(bool, tag = "2")]
pub skip_pre_flight: bool,
#[prost(bool, optional, tag = "3")]
pub front_running_protection: ::core::option::Option<bool>,
#[prost(bool, optional, tag = "8")]
pub experimental_front_running_protection: ::core::option::Option<bool>,
#[prost(bool, optional, tag = "9")]
pub snipe_transaction: ::core::option::Option<bool>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct PostSubmitRequestEntry {
#[prost(message, optional, tag = "1")]
pub transaction: ::core::option::Option<TransactionMessage>,
#[prost(bool, tag = "2")]
pub skip_pre_flight: bool,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct PostSubmitBatchRequest {
#[prost(message, repeated, tag = "1")]
pub entries: ::prost::alloc::vec::Vec<PostSubmitRequestEntry>,
#[prost(enumeration = "SubmitStrategy", tag = "2")]
pub submit_strategy: i32,
#[prost(bool, optional, tag = "3")]
pub use_bundle: ::core::option::Option<bool>,
#[prost(bool, optional, tag = "4")]
pub front_running_protection: ::core::option::Option<bool>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct PostSubmitBatchResponseEntry {
#[prost(string, tag = "1")]
pub signature: ::prost::alloc::string::String,
#[prost(string, tag = "2")]
pub error: ::prost::alloc::string::String,
#[prost(bool, tag = "3")]
pub submitted: bool,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct PostSubmitBatchResponse {
#[prost(message, repeated, tag = "1")]
pub transactions: ::prost::alloc::vec::Vec<PostSubmitBatchResponseEntry>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct PostSubmitResponse {
#[prost(string, tag = "1")]
pub signature: ::prost::alloc::string::String,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct TransactionMessage {
#[prost(string, tag = "1")]
pub content: ::prost::alloc::string::String,
#[prost(bool, tag = "2")]
pub is_cleanup: bool,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct TransactionMessageV2 {
#[prost(string, tag = "1")]
pub content: ::prost::alloc::string::String,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
#[repr(i32)]
pub enum SubmitStrategy {
PUknown = 0,
PSubmitAll = 1,
PAbortOnFirstError = 2,
PWaitForConfirmation = 3,
}
impl SubmitStrategy {
/// String value of the enum field names used in the ProtoBuf definition.
///
/// The values are not transformed in any way and thus are considered stable
/// (if the ProtoBuf definition does not change) and safe for programmatic use.
pub fn as_str_name(&self) -> &'static str {
match self {
Self::PUknown => "P_UKNOWN",
Self::PSubmitAll => "P_SUBMIT_ALL",
Self::PAbortOnFirstError => "P_ABORT_ON_FIRST_ERROR",
Self::PWaitForConfirmation => "P_WAIT_FOR_CONFIRMATION",
}
}
/// Creates an enum from field names used in the ProtoBuf definition.
pub fn from_str_name(value: &str) -> ::core::option::Option<Self> {
match value {
"P_UKNOWN" => Some(Self::PUknown),
"P_SUBMIT_ALL" => Some(Self::PSubmitAll),
"P_ABORT_ON_FIRST_ERROR" => Some(Self::PAbortOnFirstError),
"P_WAIT_FOR_CONFIRMATION" => Some(Self::PWaitForConfirmation),
_ => None,
}
}
}
/// Generated client implementations.
pub mod api_client {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
use tonic::codegen::http::Uri;
#[derive(Debug, Clone)]
pub struct ApiClient<T> {
inner: tonic::client::Grpc<T>,
}
impl ApiClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: std::convert::TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> ApiClient<T>
where
T: tonic::client::GrpcService<tonic::body::BoxBody>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
{
pub fn new(inner: T) -> Self {
let inner = tonic::client::Grpc::new(inner);
Self { inner }
}
pub fn with_origin(inner: T, origin: Uri) -> Self {
let inner = tonic::client::Grpc::with_origin(inner, origin);
Self { inner }
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> ApiClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
Response = http::Response<
<T as tonic::client::GrpcService<tonic::body::BoxBody>>::ResponseBody,
>,
>,
<T as tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
>>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
{
ApiClient::new(InterceptedService::new(inner, interceptor))
}
/// Compress requests with the given encoding.
///
/// This requires the server to support it otherwise it might respond with an
/// error.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.send_compressed(encoding);
self
}
/// Enable decompressing responses.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.accept_compressed(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_decoding_message_size(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_encoding_message_size(limit);
self
}
pub async fn post_submit_v2(
&mut self,
request: impl tonic::IntoRequest<super::PostSubmitRequest>,
) -> std::result::Result<
tonic::Response<super::PostSubmitResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static("/api.Api/PostSubmitV2");
let mut req = request.into_request();
req.extensions_mut().insert(GrpcMethod::new("api.Api", "PostSubmitV2"));
self.inner.unary(req, path, codec).await
}
pub async fn post_submit_batch_v2(
&mut self,
request: impl tonic::IntoRequest<super::PostSubmitBatchRequest>,
) -> std::result::Result<
tonic::Response<super::PostSubmitBatchResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/api.Api/PostSubmitBatchV2",
);
let mut req = request.into_request();
req.extensions_mut().insert(GrpcMethod::new("api.Api", "PostSubmitBatchV2"));
self.inner.unary(req, path, codec).await
}
}
}
/// Generated server implementations.
pub mod api_server {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
/// Generated trait containing gRPC methods that should be implemented for use with ApiServer.
#[async_trait]
pub trait Api: std::marker::Send + std::marker::Sync + 'static {
async fn post_submit_v2(
&self,
request: tonic::Request<super::PostSubmitRequest>,
) -> std::result::Result<
tonic::Response<super::PostSubmitResponse>,
tonic::Status,
>;
async fn post_submit_batch_v2(
&self,
request: tonic::Request<super::PostSubmitBatchRequest>,
) -> std::result::Result<
tonic::Response<super::PostSubmitBatchResponse>,
tonic::Status,
>;
}
#[derive(Debug)]
pub struct ApiServer<T> {
inner: Arc<T>,
accept_compression_encodings: EnabledCompressionEncodings,
send_compression_encodings: EnabledCompressionEncodings,
max_decoding_message_size: Option<usize>,
max_encoding_message_size: Option<usize>,
}
impl<T> ApiServer<T> {
pub fn new(inner: T) -> Self {
Self::from_arc(Arc::new(inner))
}
pub fn from_arc(inner: Arc<T>) -> Self {
Self {
inner,
accept_compression_encodings: Default::default(),
send_compression_encodings: Default::default(),
max_decoding_message_size: None,
max_encoding_message_size: None,
}
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> InterceptedService<Self, F>
where
F: tonic::service::Interceptor,
{
InterceptedService::new(Self::new(inner), interceptor)
}
/// Enable decompressing requests with the given encoding.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.accept_compression_encodings.enable(encoding);
self
}
/// Compress responses with the given encoding, if the client supports it.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.send_compression_encodings.enable(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.max_decoding_message_size = Some(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.max_encoding_message_size = Some(limit);
self
}
}
impl<T, B> tonic::codegen::Service<http::Request<B>> for ApiServer<T>
where
T: Api,
B: Body + std::marker::Send + 'static,
B::Error: Into<StdError> + std::marker::Send + 'static,
{
type Response = http::Response<tonic::body::BoxBody>;
type Error = std::convert::Infallible;
type Future = BoxFuture<Self::Response, Self::Error>;
fn poll_ready(
&mut self,
_cx: &mut Context<'_>,
) -> Poll<std::result::Result<(), Self::Error>> {
Poll::Ready(Ok(()))
}
fn call(&mut self, req: http::Request<B>) -> Self::Future {
match req.uri().path() {
"/api.Api/PostSubmitV2" => {
#[allow(non_camel_case_types)]
struct PostSubmitV2Svc<T: Api>(pub Arc<T>);
impl<T: Api> tonic::server::UnaryService<super::PostSubmitRequest>
for PostSubmitV2Svc<T> {
type Response = super::PostSubmitResponse;
type Future = BoxFuture<
tonic::Response<Self::Response>,
tonic::Status,
>;
fn call(
&mut self,
request: tonic::Request<super::PostSubmitRequest>,
) -> Self::Future {
let inner = Arc::clone(&self.0);
let fut = async move {
<T as Api>::post_submit_v2(&inner, request).await
};
Box::pin(fut)
}
}
let accept_compression_encodings = self.accept_compression_encodings;
let send_compression_encodings = self.send_compression_encodings;
let max_decoding_message_size = self.max_decoding_message_size;
let max_encoding_message_size = self.max_encoding_message_size;
let inner = self.inner.clone();
let fut = async move {
let method = PostSubmitV2Svc(inner);
let codec = tonic::codec::ProstCodec::default();
let mut grpc = tonic::server::Grpc::new(codec)
.apply_compression_config(
accept_compression_encodings,
send_compression_encodings,
)
.apply_max_message_size_config(
max_decoding_message_size,
max_encoding_message_size,
);
let res = grpc.unary(method, req).await;
Ok(res)
};
Box::pin(fut)
}
"/api.Api/PostSubmitBatchV2" => {
#[allow(non_camel_case_types)]
struct PostSubmitBatchV2Svc<T: Api>(pub Arc<T>);
impl<
T: Api,
> tonic::server::UnaryService<super::PostSubmitBatchRequest>
for PostSubmitBatchV2Svc<T> {
type Response = super::PostSubmitBatchResponse;
type Future = BoxFuture<
tonic::Response<Self::Response>,
tonic::Status,
>;
fn call(
&mut self,
request: tonic::Request<super::PostSubmitBatchRequest>,
) -> Self::Future {
let inner = Arc::clone(&self.0);
let fut = async move {
<T as Api>::post_submit_batch_v2(&inner, request).await
};
Box::pin(fut)
}
}
let accept_compression_encodings = self.accept_compression_encodings;
let send_compression_encodings = self.send_compression_encodings;
let max_decoding_message_size = self.max_decoding_message_size;
let max_encoding_message_size = self.max_encoding_message_size;
let inner = self.inner.clone();
let fut = async move {
let method = PostSubmitBatchV2Svc(inner);
let codec = tonic::codec::ProstCodec::default();
let mut grpc = tonic::server::Grpc::new(codec)
.apply_compression_config(
accept_compression_encodings,
send_compression_encodings,
)
.apply_max_message_size_config(
max_decoding_message_size,
max_encoding_message_size,
);
let res = grpc.unary(method, req).await;
Ok(res)
};
Box::pin(fut)
}
_ => {
Box::pin(async move {
let mut response = http::Response::new(empty_body());
let headers = response.headers_mut();
headers
.insert(
tonic::Status::GRPC_STATUS,
(tonic::Code::Unimplemented as i32).into(),
);
headers
.insert(
http::header::CONTENT_TYPE,
tonic::metadata::GRPC_CONTENT_TYPE,
);
Ok(response)
})
}
}
}
}
impl<T> Clone for ApiServer<T> {
fn clone(&self) -> Self {
let inner = self.inner.clone();
Self {
inner,
accept_compression_encodings: self.accept_compression_encodings,
send_compression_encodings: self.send_compression_encodings,
max_decoding_message_size: self.max_decoding_message_size,
max_encoding_message_size: self.max_encoding_message_size,
}
}
}
/// Generated gRPC service name
pub const SERVICE_NAME: &str = "api.Api";
impl<T> tonic::server::NamedService for ApiServer<T> {
const NAME: &'static str = SERVICE_NAME;
}
}
+41
View File
@@ -0,0 +1,41 @@
// This file is @generated by prost-build.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct PacketBatch {
#[prost(message, repeated, tag = "1")]
pub packets: ::prost::alloc::vec::Vec<Packet>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Packet {
#[prost(bytes = "vec", tag = "1")]
pub data: ::prost::alloc::vec::Vec<u8>,
#[prost(message, optional, tag = "2")]
pub meta: ::core::option::Option<Meta>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Meta {
#[prost(uint64, tag = "1")]
pub size: u64,
#[prost(string, tag = "2")]
pub addr: ::prost::alloc::string::String,
#[prost(uint32, tag = "3")]
pub port: u32,
#[prost(message, optional, tag = "4")]
pub flags: ::core::option::Option<PacketFlags>,
#[prost(uint64, tag = "5")]
pub sender_stake: u64,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct PacketFlags {
#[prost(bool, tag = "1")]
pub discard: bool,
#[prost(bool, tag = "2")]
pub forwarded: bool,
#[prost(bool, tag = "3")]
pub repair: bool,
#[prost(bool, tag = "4")]
pub simple_vote_tx: bool,
#[prost(bool, tag = "5")]
pub tracer_packet: bool,
#[prost(bool, tag = "6")]
pub from_staked_node: bool,
}
+178
View File
@@ -0,0 +1,178 @@
// This file is @generated by prost-build.
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct GetTpuConfigsRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct GetTpuConfigsResponse {
#[prost(message, optional, tag = "1")]
pub tpu: ::core::option::Option<super::shared::Socket>,
#[prost(message, optional, tag = "2")]
pub tpu_forward: ::core::option::Option<super::shared::Socket>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct SubscribePacketsRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct SubscribePacketsResponse {
#[prost(message, optional, tag = "1")]
pub header: ::core::option::Option<super::shared::Header>,
#[prost(oneof = "subscribe_packets_response::Msg", tags = "2, 3")]
pub msg: ::core::option::Option<subscribe_packets_response::Msg>,
}
/// Nested message and enum types in `SubscribePacketsResponse`.
pub mod subscribe_packets_response {
#[derive(Clone, PartialEq, ::prost::Oneof)]
pub enum Msg {
#[prost(message, tag = "2")]
Heartbeat(super::super::shared::Heartbeat),
#[prost(message, tag = "3")]
Batch(super::super::packet::PacketBatch),
}
}
/// Generated client implementations.
pub mod relayer_client {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
use tonic::codegen::http::Uri;
/// / Relayers offer a TPU and TPU forward proxy for Solana validators.
/// / Validators can connect and fetch the TPU configuration for the relayer and start to advertise the
/// / relayer's information in gossip.
/// / They can also subscribe to packets which arrived on the TPU ports at the relayer
#[derive(Debug, Clone)]
pub struct RelayerClient<T> {
inner: tonic::client::Grpc<T>,
}
impl RelayerClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> RelayerClient<T>
where
T: tonic::client::GrpcService<tonic::body::BoxBody>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
{
pub fn new(inner: T) -> Self {
let inner = tonic::client::Grpc::new(inner);
Self { inner }
}
pub fn with_origin(inner: T, origin: Uri) -> Self {
let inner = tonic::client::Grpc::with_origin(inner, origin);
Self { inner }
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> RelayerClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
Response = http::Response<
<T as tonic::client::GrpcService<tonic::body::BoxBody>>::ResponseBody,
>,
>,
<T as tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
>>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
{
RelayerClient::new(InterceptedService::new(inner, interceptor))
}
/// Compress requests with the given encoding.
///
/// This requires the server to support it otherwise it might respond with an
/// error.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.send_compressed(encoding);
self
}
/// Enable decompressing responses.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.accept_compressed(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_decoding_message_size(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_encoding_message_size(limit);
self
}
/// The relayer has TPU and TPU forward sockets that validators can leverage.
/// A validator can fetch this config and change its TPU and TPU forward port in gossip.
pub async fn get_tpu_configs(
&mut self,
request: impl tonic::IntoRequest<super::GetTpuConfigsRequest>,
) -> std::result::Result<
tonic::Response<super::GetTpuConfigsResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/relayer.Relayer/GetTpuConfigs",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("relayer.Relayer", "GetTpuConfigs"));
self.inner.unary(req, path, codec).await
}
/// Validators can subscribe to packets from the relayer and receive a multiplexed signal that contains a mixture
/// of packets and heartbeats
pub async fn subscribe_packets(
&mut self,
request: impl tonic::IntoRequest<super::SubscribePacketsRequest>,
) -> std::result::Result<
tonic::Response<tonic::codec::Streaming<super::SubscribePacketsResponse>>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/relayer.Relayer/SubscribePackets",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("relayer.Relayer", "SubscribePackets"));
self.inner.server_streaming(req, path, codec).await
}
}
}
+363
View File
@@ -0,0 +1,363 @@
// This file is @generated by prost-build.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct SlotList {
#[prost(uint64, repeated, tag = "1")]
pub slots: ::prost::alloc::vec::Vec<u64>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct ConnectedLeadersResponse {
/// Mapping of validator pubkey to leader slots for the current epoch.
#[prost(map = "string, message", tag = "1")]
pub connected_validators: ::std::collections::HashMap<
::prost::alloc::string::String,
SlotList,
>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct SendBundleRequest {
#[prost(message, optional, tag = "1")]
pub bundle: ::core::option::Option<super::bundle::Bundle>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct SendBundleResponse {
/// server uuid for the bundle
#[prost(string, tag = "1")]
pub uuid: ::prost::alloc::string::String,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct NextScheduledLeaderRequest {
/// Defaults to the currently connected region if no region provided.
#[prost(string, repeated, tag = "1")]
pub regions: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct NextScheduledLeaderResponse {
/// the current slot the backend is on
#[prost(uint64, tag = "1")]
pub current_slot: u64,
/// the slot of the next leader
#[prost(uint64, tag = "2")]
pub next_leader_slot: u64,
/// the identity pubkey (base58) of the next leader
#[prost(string, tag = "3")]
pub next_leader_identity: ::prost::alloc::string::String,
/// the block engine region of the next leader
#[prost(string, tag = "4")]
pub next_leader_region: ::prost::alloc::string::String,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct ConnectedLeadersRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct ConnectedLeadersRegionedRequest {
/// Defaults to the currently connected region if no region provided.
#[prost(string, repeated, tag = "1")]
pub regions: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct ConnectedLeadersRegionedResponse {
#[prost(map = "string, message", tag = "1")]
pub connected_validators: ::std::collections::HashMap<
::prost::alloc::string::String,
ConnectedLeadersResponse,
>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct GetTipAccountsRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct GetTipAccountsResponse {
#[prost(string, repeated, tag = "1")]
pub accounts: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct SubscribeBundleResultsRequest {}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct GetRegionsRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct GetRegionsResponse {
/// The region the client is currently connected to
#[prost(string, tag = "1")]
pub current_region: ::prost::alloc::string::String,
/// Regions that are online and ready for connections
/// All regions: <https://jito-labs.gitbook.io/mev/systems/connecting/mainnet>
#[prost(string, repeated, tag = "2")]
pub available_regions: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
}
/// Generated client implementations.
pub mod searcher_service_client {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
use tonic::codegen::http::Uri;
#[derive(Debug, Clone)]
pub struct SearcherServiceClient<T> {
inner: tonic::client::Grpc<T>,
}
impl SearcherServiceClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> SearcherServiceClient<T>
where
T: tonic::client::GrpcService<tonic::body::BoxBody>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
{
pub fn new(inner: T) -> Self {
let inner = tonic::client::Grpc::new(inner);
Self { inner }
}
pub fn with_origin(inner: T, origin: Uri) -> Self {
let inner = tonic::client::Grpc::with_origin(inner, origin);
Self { inner }
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> SearcherServiceClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
Response = http::Response<
<T as tonic::client::GrpcService<tonic::body::BoxBody>>::ResponseBody,
>,
>,
<T as tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
>>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
{
SearcherServiceClient::new(InterceptedService::new(inner, interceptor))
}
/// Compress requests with the given encoding.
///
/// This requires the server to support it otherwise it might respond with an
/// error.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.send_compressed(encoding);
self
}
/// Enable decompressing responses.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.accept_compressed(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_decoding_message_size(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_encoding_message_size(limit);
self
}
/// Searchers can invoke this endpoint to subscribe to their respective bundle results.
/// A success result would indicate the bundle won its state auction and was submitted to the validator.
pub async fn subscribe_bundle_results(
&mut self,
request: impl tonic::IntoRequest<super::SubscribeBundleResultsRequest>,
) -> std::result::Result<
tonic::Response<tonic::codec::Streaming<super::super::bundle::BundleResult>>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/searcher.SearcherService/SubscribeBundleResults",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new("searcher.SearcherService", "SubscribeBundleResults"),
);
self.inner.server_streaming(req, path, codec).await
}
pub async fn send_bundle(
&mut self,
request: impl tonic::IntoRequest<super::SendBundleRequest>,
) -> std::result::Result<
tonic::Response<super::SendBundleResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/searcher.SearcherService/SendBundle",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("searcher.SearcherService", "SendBundle"));
self.inner.unary(req, path, codec).await
}
/// Returns the next scheduled leader connected to the block engine.
pub async fn get_next_scheduled_leader(
&mut self,
request: impl tonic::IntoRequest<super::NextScheduledLeaderRequest>,
) -> std::result::Result<
tonic::Response<super::NextScheduledLeaderResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/searcher.SearcherService/GetNextScheduledLeader",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new("searcher.SearcherService", "GetNextScheduledLeader"),
);
self.inner.unary(req, path, codec).await
}
/// Returns leader slots for connected jito validators during the current epoch. Only returns data for this region.
pub async fn get_connected_leaders(
&mut self,
request: impl tonic::IntoRequest<super::ConnectedLeadersRequest>,
) -> std::result::Result<
tonic::Response<super::ConnectedLeadersResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/searcher.SearcherService/GetConnectedLeaders",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new("searcher.SearcherService", "GetConnectedLeaders"),
);
self.inner.unary(req, path, codec).await
}
/// Returns leader slots for connected jito validators during the current epoch.
pub async fn get_connected_leaders_regioned(
&mut self,
request: impl tonic::IntoRequest<super::ConnectedLeadersRegionedRequest>,
) -> std::result::Result<
tonic::Response<super::ConnectedLeadersRegionedResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/searcher.SearcherService/GetConnectedLeadersRegioned",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new(
"searcher.SearcherService",
"GetConnectedLeadersRegioned",
),
);
self.inner.unary(req, path, codec).await
}
/// Returns the tip accounts searchers shall transfer funds to for the leader to claim.
pub async fn get_tip_accounts(
&mut self,
request: impl tonic::IntoRequest<super::GetTipAccountsRequest>,
) -> std::result::Result<
tonic::Response<super::GetTipAccountsResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/searcher.SearcherService/GetTipAccounts",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("searcher.SearcherService", "GetTipAccounts"));
self.inner.unary(req, path, codec).await
}
/// Returns region the client is directly connected to, along with all available regions
pub async fn get_regions(
&mut self,
request: impl tonic::IntoRequest<super::GetRegionsRequest>,
) -> std::result::Result<
tonic::Response<super::GetRegionsResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/searcher.SearcherService/GetRegions",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("searcher.SearcherService", "GetRegions"));
self.inner.unary(req, path, codec).await
}
}
}
+131
View File
@@ -0,0 +1,131 @@
use std::{
sync::Arc,
time::{Duration, Instant},
};
use crate::protos::{
bundle::{
Bundle, BundleResult,
},
convert::proto_packet_from_versioned_tx,
searcher::{
searcher_service_client::SearcherServiceClient, SendBundleRequest, SubscribeBundleResultsRequest,
},
};
use solana_sdk::{
signature::Signature,
transaction::VersionedTransaction,
};
use thiserror::Error;
use tokio::sync::Mutex;
use tonic::{
transport::{self, Channel, Endpoint}, Status
};
use yellowstone_grpc_client::ClientTlsConfig;
use crate::swqos::common::poll_transaction_confirmation;
use crate::common::SolanaRpcClient;
use crate::swqos::TradeType;
#[derive(Debug, Error)]
pub enum BlockEngineConnectionError {
#[error("transport error {0}")]
TransportError(#[from] transport::Error),
#[error("client error {0}")]
ClientError(#[from] Status),
}
#[derive(Debug, Error)]
pub enum BundleRejectionError {
#[error("bundle lost state auction, auction: {0}, tip {1} lamports")]
StateAuctionBidRejected(String, u64),
#[error("bundle won state auction but failed global auction, auction {0}, tip {1} lamports")]
WinningBatchBidRejected(String, u64),
#[error("bundle simulation failure on tx {0}, message: {1:?}")]
SimulationFailure(String, Option<String>),
#[error("internal error {0}")]
InternalError(String),
}
pub type BlockEngineConnectionResult<T> = Result<T, BlockEngineConnectionError>;
pub async fn get_searcher_client_no_auth(
block_engine_url: &str,
) -> BlockEngineConnectionResult<SearcherServiceClient<Channel>> {
let searcher_channel = create_grpc_channel(block_engine_url).await?;
let searcher_client = SearcherServiceClient::new(searcher_channel);
Ok(searcher_client)
}
pub async fn create_grpc_channel(url: &str) -> BlockEngineConnectionResult<Channel> {
let mut endpoint = Endpoint::from_shared(url.to_string()).expect("invalid url");
if url.starts_with("https") {
endpoint = endpoint.tls_config(ClientTlsConfig::new().with_native_roots())?;
}
endpoint = endpoint.tcp_nodelay(true);
endpoint = endpoint.tcp_keepalive(Some(Duration::from_secs(10)));
endpoint = endpoint.connect_timeout(Duration::from_secs(20));
endpoint = endpoint.http2_keep_alive_interval(Duration::from_secs(10));
Ok(endpoint.connect().await?)
}
pub async fn subscribe_bundle_results(
searcher_client: Arc<Mutex<SearcherServiceClient<Channel>>>,
request: impl tonic::IntoRequest<SubscribeBundleResultsRequest>,
) -> std::result::Result<
tonic::Response<tonic::codec::Streaming<BundleResult>>,
tonic::Status,
> {
let mut searcher = searcher_client.lock().await;
searcher.subscribe_bundle_results(request).await
}
pub async fn send_bundle_with_confirmation(
rpc: Arc<SolanaRpcClient>,
trade_type: TradeType,
transactions: &Vec<VersionedTransaction>,
searcher_client: Arc<Mutex<SearcherServiceClient<Channel>>>,
) -> Result<Vec<Signature>, anyhow::Error> {
let start_time = Instant::now();
let signatures = send_bundle_no_wait(transactions, searcher_client).await?;
println!(" Jito{}提交: {:?}", trade_type, start_time.elapsed());
let start_time: Instant = Instant::now();
for signature in signatures.clone() {
match poll_transaction_confirmation(&rpc, signature).await {
Ok(_) => continue,
Err(_) => continue,
}
}
println!(" Jito{}确认: {:?}", trade_type, start_time.elapsed());
Ok(signatures)
}
pub async fn send_bundle_no_wait(
transactions: &Vec<VersionedTransaction>,
searcher_client: Arc<Mutex<SearcherServiceClient<Channel>>>,
) -> Result<Vec<Signature>, anyhow::Error> {
let mut packets = vec![];
let mut signatures = vec![];
for transaction in transactions {
let packet = proto_packet_from_versioned_tx(transaction);
packets.push(packet);
signatures.push(transaction.signatures[0]);
}
let mut searcher = searcher_client.lock().await;
searcher
.send_bundle(SendBundleRequest {
bundle: Some(Bundle {
header: None,
packets,
}),
})
.await?;
Ok(signatures)
}
+18
View File
@@ -0,0 +1,18 @@
// This file is @generated by prost-build.
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct Header {
#[prost(message, optional, tag = "1")]
pub ts: ::core::option::Option<::prost_types::Timestamp>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct Heartbeat {
#[prost(uint64, tag = "1")]
pub count: u64,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Socket {
#[prost(string, tag = "1")]
pub ip: ::prost::alloc::string::String,
#[prost(int64, tag = "2")]
pub port: i64,
}
+279
View File
@@ -0,0 +1,279 @@
// This file is @generated by prost-build.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Heartbeat {
/// don't trust IP:PORT from tcp header since it can be tampered over the wire
/// `socket.ip` must match incoming packet's ip. this prevents spamming an unwitting destination
#[prost(message, optional, tag = "1")]
pub socket: ::core::option::Option<super::shared::Socket>,
/// regions for shredstream proxy to receive shreds from
/// list of valid regions: <https://docs.jito.wtf/lowlatencytxnsend/#api>
#[prost(string, repeated, tag = "2")]
pub regions: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct HeartbeatResponse {
/// client must respond within `ttl_ms` to keep stream alive
#[prost(uint32, tag = "1")]
pub ttl_ms: u32,
}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct TraceShred {
/// source region, one of: <https://docs.jito.wtf/lowlatencytxnsend/#api>
#[prost(string, tag = "1")]
pub region: ::prost::alloc::string::String,
/// timestamp of creation
#[prost(message, optional, tag = "2")]
pub created_at: ::core::option::Option<::prost_types::Timestamp>,
/// monotonically increases, resets upon service restart
#[prost(uint32, tag = "3")]
pub seq_num: u32,
}
/// tbd: we may want to add filters here
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct SubscribeEntriesRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct Entry {
/// the slot that the entry is from
#[prost(uint64, tag = "1")]
pub slot: u64,
/// Serialized bytes of Vec<Entry>: <https://docs.rs/solana-entry/latest/solana_entry/entry/struct.Entry.html>
#[prost(bytes = "vec", tag = "2")]
pub entries: ::prost::alloc::vec::Vec<u8>,
}
/// Generated client implementations.
pub mod shredstream_client {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
use tonic::codegen::http::Uri;
#[derive(Debug, Clone)]
pub struct ShredstreamClient<T> {
inner: tonic::client::Grpc<T>,
}
impl ShredstreamClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> ShredstreamClient<T>
where
T: tonic::client::GrpcService<tonic::body::BoxBody>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
{
pub fn new(inner: T) -> Self {
let inner = tonic::client::Grpc::new(inner);
Self { inner }
}
pub fn with_origin(inner: T, origin: Uri) -> Self {
let inner = tonic::client::Grpc::with_origin(inner, origin);
Self { inner }
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> ShredstreamClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
Response = http::Response<
<T as tonic::client::GrpcService<tonic::body::BoxBody>>::ResponseBody,
>,
>,
<T as tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
>>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
{
ShredstreamClient::new(InterceptedService::new(inner, interceptor))
}
/// Compress requests with the given encoding.
///
/// This requires the server to support it otherwise it might respond with an
/// error.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.send_compressed(encoding);
self
}
/// Enable decompressing responses.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.accept_compressed(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_decoding_message_size(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_encoding_message_size(limit);
self
}
/// RPC endpoint to send heartbeats to keep shreds flowing
pub async fn send_heartbeat(
&mut self,
request: impl tonic::IntoRequest<super::Heartbeat>,
) -> std::result::Result<
tonic::Response<super::HeartbeatResponse>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/shredstream.Shredstream/SendHeartbeat",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("shredstream.Shredstream", "SendHeartbeat"));
self.inner.unary(req, path, codec).await
}
}
}
/// Generated client implementations.
pub mod shredstream_proxy_client {
#![allow(
unused_variables,
dead_code,
missing_docs,
clippy::wildcard_imports,
clippy::let_unit_value,
)]
use tonic::codegen::*;
use tonic::codegen::http::Uri;
#[derive(Debug, Clone)]
pub struct ShredstreamProxyClient<T> {
inner: tonic::client::Grpc<T>,
}
impl ShredstreamProxyClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> ShredstreamProxyClient<T>
where
T: tonic::client::GrpcService<tonic::body::BoxBody>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + std::marker::Send,
{
pub fn new(inner: T) -> Self {
let inner = tonic::client::Grpc::new(inner);
Self { inner }
}
pub fn with_origin(inner: T, origin: Uri) -> Self {
let inner = tonic::client::Grpc::with_origin(inner, origin);
Self { inner }
}
pub fn with_interceptor<F>(
inner: T,
interceptor: F,
) -> ShredstreamProxyClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
Response = http::Response<
<T as tonic::client::GrpcService<tonic::body::BoxBody>>::ResponseBody,
>,
>,
<T as tonic::codegen::Service<
http::Request<tonic::body::BoxBody>,
>>::Error: Into<StdError> + std::marker::Send + std::marker::Sync,
{
ShredstreamProxyClient::new(InterceptedService::new(inner, interceptor))
}
/// Compress requests with the given encoding.
///
/// This requires the server to support it otherwise it might respond with an
/// error.
#[must_use]
pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.send_compressed(encoding);
self
}
/// Enable decompressing responses.
#[must_use]
pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self {
self.inner = self.inner.accept_compressed(encoding);
self
}
/// Limits the maximum size of a decoded message.
///
/// Default: `4MB`
#[must_use]
pub fn max_decoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_decoding_message_size(limit);
self
}
/// Limits the maximum size of an encoded message.
///
/// Default: `usize::MAX`
#[must_use]
pub fn max_encoding_message_size(mut self, limit: usize) -> Self {
self.inner = self.inner.max_encoding_message_size(limit);
self
}
pub async fn subscribe_entries(
&mut self,
request: impl tonic::IntoRequest<super::SubscribeEntriesRequest>,
) -> std::result::Result<
tonic::Response<tonic::codec::Streaming<super::Entry>>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic::codec::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/shredstream.ShredstreamProxy/SubscribeEntries",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new("shredstream.ShredstreamProxy", "SubscribeEntries"),
);
self.inner.server_streaming(req, path, codec).await
}
}
}
+167
View File
@@ -0,0 +1,167 @@
use std::{
sync::{Arc, RwLock},
time::{Duration, SystemTime},
};
use crate::protos::auth::{
auth_service_client::AuthServiceClient, GenerateAuthChallengeRequest,
GenerateAuthTokensRequest, RefreshAccessTokenRequest, Role, Token,
};
use prost_types::Timestamp;
use solana_metrics::datapoint_info;
use solana_sdk::signature::{Keypair, Signer};
use tokio::{task::JoinHandle, time::sleep};
use tonic::{service::Interceptor, transport::Channel, Request, Status};
use super::searcher_client::BlockEngineConnectionResult;
const AUTHORIZATION_HEADER: &str = "authorization";
const BEARER: &str = "Bearer ";
/// Adds the token to each requests' authorization header.
/// Manages refreshing the token in a separate thread.
#[derive(Clone)]
pub struct ClientInterceptor {
/// The token added to each request header.
bearer_token: Arc<RwLock<String>>,
}
impl ClientInterceptor {
pub async fn new(
mut auth_service_client: AuthServiceClient<Channel>,
keypair: &Arc<Keypair>,
role: Role,
) -> BlockEngineConnectionResult<Self> {
let (access_token, refresh_token) =
Self::auth(&mut auth_service_client, keypair, role).await?;
let bearer_token = Arc::new(RwLock::new(access_token.value.clone()));
let _refresh_token_thread = Self::spawn_token_refresh_thread(
auth_service_client,
bearer_token.clone(),
refresh_token,
access_token.expires_at_utc.unwrap(),
keypair.clone(),
role,
);
Ok(Self { bearer_token })
}
async fn auth(
auth_service_client: &mut AuthServiceClient<Channel>,
keypair: &Keypair,
role: Role,
) -> BlockEngineConnectionResult<(Token, Token)> {
let challenge_resp = auth_service_client
.generate_auth_challenge(GenerateAuthChallengeRequest {
role: role as i32,
pubkey: keypair.pubkey().as_ref().to_vec(),
})
.await?
.into_inner();
let challenge = format!("{}-{}", keypair.pubkey(), challenge_resp.challenge);
let signed_challenge = keypair.sign_message(challenge.as_bytes()).as_ref().to_vec();
let tokens = auth_service_client
.generate_auth_tokens(GenerateAuthTokensRequest {
challenge,
client_pubkey: keypair.pubkey().as_ref().to_vec(),
signed_challenge,
})
.await?
.into_inner();
Ok((tokens.access_token.unwrap(), tokens.refresh_token.unwrap()))
}
fn spawn_token_refresh_thread(
mut auth_service_client: AuthServiceClient<Channel>,
bearer_token: Arc<RwLock<String>>,
refresh_token: Token,
access_token_expiration: Timestamp,
keypair: Arc<Keypair>,
role: Role,
) -> JoinHandle<BlockEngineConnectionResult<()>> {
tokio::spawn(async move {
let mut refresh_token = refresh_token;
let mut access_token_expiration = access_token_expiration;
loop {
let access_token_ttl = SystemTime::try_from(access_token_expiration.clone())
.unwrap()
.duration_since(SystemTime::now())
.unwrap_or_else(|_| Duration::from_secs(0));
let refresh_token_ttl =
SystemTime::try_from(refresh_token.expires_at_utc.as_ref().unwrap().clone())
.unwrap()
.duration_since(SystemTime::now())
.unwrap_or_else(|_| Duration::from_secs(0));
let does_access_token_expire_soon = access_token_ttl < Duration::from_secs(5 * 60);
let does_refresh_token_expire_soon =
refresh_token_ttl < Duration::from_secs(5 * 60);
match (
does_refresh_token_expire_soon,
does_access_token_expire_soon,
) {
// re-run entire auth workflow is refresh token expiring soon
(true, _) => {
let is_error = {
if let Ok((new_access_token, new_refresh_token)) =
Self::auth(&mut auth_service_client, &keypair, role).await
{
*bearer_token.write().unwrap() = new_access_token.value.clone();
access_token_expiration = new_access_token.expires_at_utc.unwrap();
refresh_token = new_refresh_token;
false
} else {
true
}
};
datapoint_info!("searcher-full-auth", ("is_error", is_error, bool));
}
// re-up the access token if it expires soon
(_, true) => {
let is_error = {
if let Ok(refresh_resp) = auth_service_client
.refresh_access_token(RefreshAccessTokenRequest {
refresh_token: refresh_token.value.clone(),
})
.await
{
let access_token = refresh_resp.into_inner().access_token.unwrap();
*bearer_token.write().unwrap() = access_token.value.clone();
access_token_expiration = access_token.expires_at_utc.unwrap();
false
} else {
true
}
};
datapoint_info!("searcher-refresh-auth", ("is_error", is_error, bool));
}
_ => {
sleep(Duration::from_secs(60)).await;
}
}
}
})
}
}
impl Interceptor for ClientInterceptor {
fn call(&mut self, mut request: Request<()>) -> Result<Request<()>, Status> {
let l_token = self.bearer_token.read().unwrap();
if !l_token.is_empty() {
request.metadata_mut().insert(
AUTHORIZATION_HEADER,
format!("{BEARER}{l_token}").parse().unwrap(),
);
}
Ok(request)
}
}
+13
View File
@@ -0,0 +1,13 @@
// This file is @generated by prost-build.
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct TraceShred {
/// source region, one of: <https://jito-labs.gitbook.io/mev/systems/connecting/mainnet>
#[prost(string, tag = "1")]
pub region: ::prost::alloc::string::String,
/// timestamp of creation
#[prost(message, optional, tag = "2")]
pub created_at: ::core::option::Option<::prost_types::Timestamp>,
/// monotonically increases, resets upon service restart
#[prost(uint32, tag = "3")]
pub seq_num: u32,
}