v0.3.14-pre.006
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs
|
||||
// version: 31
|
||||
// version: 32
|
||||
|
||||
use sha2::Digest; // rust-rules: trait-import
|
||||
|
||||
@@ -36,6 +36,42 @@ const RAW_TRANSACTION_INGEST_YELLOWSTONE_FILTER_FINGERPRINT_DOMAIN: &[u8] = b"ks
|
||||
const RAW_TRANSACTION_INGEST_YELLOWSTONE_HTTP_PROTOCOL: &str = "yellowstone_http";
|
||||
const RAW_TRANSACTION_INGEST_YELLOWSTONE_HTTP_SOURCE_KEY_DOMAIN: &[u8] = b"ksp.raw_transaction_ingest.yellowstone_http.source_key.v1\0";
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
enum RawTransactionIngestHttpDiscoveryStrategy {
|
||||
ClosedRange,
|
||||
WithLimitBoundary,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Eq, PartialEq)]
|
||||
struct RawTransactionIngestHttpDiscoveryWindow {
|
||||
produced_slots: std::vec::Vec<u64>,
|
||||
proven_end_slot: std::option::Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
struct RawTransactionIngestHttpScanCapabilities {
|
||||
get_block: bool,
|
||||
get_blocks: bool,
|
||||
get_blocks_with_limit: bool,
|
||||
get_slot: bool,
|
||||
}
|
||||
|
||||
impl RawTransactionIngestHttpScanCapabilities {
|
||||
const fn can_scan(self) -> bool {
|
||||
return self.get_block && (self.get_blocks || (self.get_blocks_with_limit && self.get_slot));
|
||||
}
|
||||
|
||||
fn strategy(self) -> ksp_core_lib::Result<RawTransactionIngestHttpDiscoveryStrategy> {
|
||||
if self.get_block && self.get_blocks {
|
||||
return std::result::Result::Ok(RawTransactionIngestHttpDiscoveryStrategy::ClosedRange);
|
||||
}
|
||||
if self.get_block && self.get_blocks_with_limit && self.get_slot {
|
||||
return std::result::Result::Ok(RawTransactionIngestHttpDiscoveryStrategy::WithLimitBoundary);
|
||||
}
|
||||
return std::result::Result::Err(crate::runtime_error("continuity.http_scan_capability_incomplete"));
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
enum RawTransactionIngestSourceFamily {
|
||||
Block,
|
||||
@@ -1340,33 +1376,40 @@ impl crate::RawTransactionIngestHttpBlockPollingSource {
|
||||
},
|
||||
};
|
||||
if next_scan_slot <= current_tip {
|
||||
let discovered = tokio::select! {
|
||||
biased;
|
||||
_ = stop_receiver.changed() => {
|
||||
break;
|
||||
}
|
||||
result = self.http_pool.get_blocks_with_limit(
|
||||
&self.polling_role,
|
||||
next_scan_slot,
|
||||
u64::from(self.max_discovered_blocks_per_cycle),
|
||||
std::option::Option::Some(&context_config),
|
||||
) => result,
|
||||
};
|
||||
let discovered = match discovered {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => {
|
||||
fault = std::option::Option::Some(source_transport_error(error.code()));
|
||||
let configured_window_slots =
|
||||
u64::from(self.max_discovered_blocks_per_cycle).min(crate::MAX_RAW_TRANSACTION_INGEST_CONTINUITY_DISCOVERY_WINDOW_SLOTS);
|
||||
let window_offset = match configured_window_slots.checked_sub(1) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => {
|
||||
fault = std::option::Option::Some(crate::runtime_error("source.http_block_polling_window_invalid"));
|
||||
break;
|
||||
},
|
||||
};
|
||||
let candidate_end = match next_scan_slot.checked_add(window_offset) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => u64::MAX,
|
||||
};
|
||||
let window_end = current_tip.min(candidate_end);
|
||||
let discovered = discover_http_block_window(
|
||||
&self.http_pool,
|
||||
&self.polling_role,
|
||||
self.network.as_str(),
|
||||
self.commitment,
|
||||
next_scan_slot,
|
||||
window_end,
|
||||
&mut stop_receiver,
|
||||
)
|
||||
.await;
|
||||
let discovery = match discovered {
|
||||
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
||||
std::result::Result::Ok(std::option::Option::None) => break,
|
||||
std::result::Result::Err(error) => {
|
||||
fault = std::option::Option::Some(error);
|
||||
break;
|
||||
},
|
||||
};
|
||||
let validated = validate_http_block_polling_discovery(next_scan_slot, discovered.as_slice());
|
||||
if let std::result::Result::Err(error) = validated {
|
||||
fault = std::option::Option::Some(error);
|
||||
break;
|
||||
}
|
||||
let discovered_count = discovered.len();
|
||||
let mut blocked_by_null = false;
|
||||
for slot in discovered {
|
||||
for slot in discovery.produced_slots {
|
||||
let observed = tokio::select! {
|
||||
biased;
|
||||
_ = stop_receiver.changed() => {
|
||||
@@ -1424,16 +1467,9 @@ impl crate::RawTransactionIngestHttpBlockPollingSource {
|
||||
fault = std::option::Option::Some(error);
|
||||
break 'source;
|
||||
}
|
||||
next_scan_slot = match slot.checked_add(1) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => {
|
||||
fault = std::option::Option::Some(crate::runtime_error("source.http_block_polling_slot_exhausted"));
|
||||
break 'source;
|
||||
},
|
||||
};
|
||||
}
|
||||
if !blocked_by_null && discovered_count < usize::from(self.max_discovered_blocks_per_cycle) && next_scan_slot <= current_tip {
|
||||
next_scan_slot = match current_tip.checked_add(1) {
|
||||
if !blocked_by_null && let std::option::Option::Some(proven_end_slot) = discovery.proven_end_slot {
|
||||
next_scan_slot = match proven_end_slot.checked_add(1) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => {
|
||||
fault = std::option::Option::Some(crate::runtime_error("source.http_block_polling_slot_exhausted"));
|
||||
@@ -2399,28 +2435,40 @@ async fn drain_live_source_tasks(
|
||||
};
|
||||
}
|
||||
|
||||
fn http_role_scan_capabilities(
|
||||
pool: &ksp_onchain_transport_lib::HttpTransportPool,
|
||||
role: &ksp_onchain_transport_lib::HttpRoleName,
|
||||
expected_cluster: &str,
|
||||
) -> ksp_core_lib::Result<RawTransactionIngestHttpScanCapabilities> {
|
||||
let get_block = match http_role_supports_rpc_method(pool, role, "getBlock", expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let get_blocks = match http_role_supports_rpc_method(pool, role, "getBlocks", expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let get_blocks_with_limit = match http_role_supports_rpc_method(pool, role, "getBlocksWithLimit", expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let get_slot = match http_role_supports_rpc_method(pool, role, "getSlot", expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
return std::result::Result::Ok(RawTransactionIngestHttpScanCapabilities { get_block, get_blocks, get_blocks_with_limit, get_slot });
|
||||
}
|
||||
|
||||
fn http_role_supports_repair_scan(
|
||||
pool: &ksp_onchain_transport_lib::HttpTransportPool,
|
||||
role: &ksp_onchain_transport_lib::HttpRoleName,
|
||||
expected_cluster: &str,
|
||||
) -> ksp_core_lib::Result<bool> {
|
||||
let has_get_block = match http_role_supports_rpc_method(pool, role, "getBlock", expected_cluster) {
|
||||
let capabilities = match http_role_scan_capabilities(pool, role, expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let has_get_blocks = match http_role_supports_rpc_method(pool, role, "getBlocks", expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let has_get_blocks_with_limit = match http_role_supports_rpc_method(pool, role, "getBlocksWithLimit", expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let has_get_slot = match http_role_supports_rpc_method(pool, role, "getSlot", expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
return std::result::Result::Ok(has_get_block && has_get_slot && (has_get_blocks || has_get_blocks_with_limit));
|
||||
return std::result::Result::Ok(capabilities.can_scan());
|
||||
}
|
||||
|
||||
fn http_role_supports_rpc_method(
|
||||
@@ -2561,6 +2609,125 @@ fn http_role_supports_request_kind(role: &ksp_onchain_transport_lib::HttpEndpoin
|
||||
return role.request_kinds().iter().any(|kind| return kind.as_str() == "*" || kind.as_str() == request_kind);
|
||||
}
|
||||
|
||||
fn http_discovery_window_slot_count(start_slot: u64, end_slot: u64) -> ksp_core_lib::Result<u64> {
|
||||
if end_slot < start_slot {
|
||||
return std::result::Result::Err(crate::runtime_error("continuity.http_discovery_range_reversed"));
|
||||
}
|
||||
let slot_count = match end_slot.checked_sub(start_slot).and_then(|value| return value.checked_add(1)) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("continuity.http_discovery_range_overflow")),
|
||||
};
|
||||
if slot_count > crate::MAX_RAW_TRANSACTION_INGEST_CONTINUITY_DISCOVERY_WINDOW_SLOTS {
|
||||
return std::result::Result::Err(crate::runtime_error("continuity.http_discovery_window_too_large"));
|
||||
}
|
||||
return std::result::Result::Ok(slot_count);
|
||||
}
|
||||
|
||||
fn validate_http_block_discovery_result(
|
||||
start_slot: u64,
|
||||
end_slot: u64,
|
||||
strategy: RawTransactionIngestHttpDiscoveryStrategy,
|
||||
current_tip: std::option::Option<u64>,
|
||||
discovered: &[u64],
|
||||
) -> ksp_core_lib::Result<RawTransactionIngestHttpDiscoveryWindow> {
|
||||
if let std::result::Result::Err(error) = http_discovery_window_slot_count(start_slot, end_slot) {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
if let std::result::Result::Err(error) = validate_http_block_polling_discovery(start_slot, discovered) {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
return match strategy {
|
||||
RawTransactionIngestHttpDiscoveryStrategy::ClosedRange => {
|
||||
if discovered.iter().any(|slot| return *slot > end_slot) {
|
||||
return std::result::Result::Err(crate::runtime_error("continuity.http_closed_range_exceeded"));
|
||||
}
|
||||
std::result::Result::Ok(RawTransactionIngestHttpDiscoveryWindow {
|
||||
produced_slots: discovered.to_vec(),
|
||||
proven_end_slot: std::option::Option::Some(end_slot),
|
||||
})
|
||||
},
|
||||
RawTransactionIngestHttpDiscoveryStrategy::WithLimitBoundary => {
|
||||
let current_tip = match current_tip {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("continuity.http_discovery_tip_missing")),
|
||||
};
|
||||
if current_tip < start_slot && discovered.is_empty() {
|
||||
return std::result::Result::Ok(RawTransactionIngestHttpDiscoveryWindow {
|
||||
produced_slots: std::vec::Vec::new(),
|
||||
proven_end_slot: std::option::Option::None,
|
||||
});
|
||||
}
|
||||
let proven_end_slot = discovered.last().copied().map(|slot| return slot.min(end_slot));
|
||||
let produced_slots = discovered.iter().copied().take_while(|slot| return *slot <= end_slot).collect();
|
||||
std::result::Result::Ok(RawTransactionIngestHttpDiscoveryWindow { produced_slots, proven_end_slot })
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
async fn discover_http_block_window(
|
||||
pool: &ksp_onchain_transport_lib::HttpTransportPool,
|
||||
role: &ksp_onchain_transport_lib::HttpRoleName,
|
||||
expected_cluster: &str,
|
||||
commitment: ksp_onchain_transport_lib::SolanaCommitment,
|
||||
start_slot: u64,
|
||||
end_slot: u64,
|
||||
stop_receiver: &mut tokio::sync::watch::Receiver<bool>,
|
||||
) -> ksp_core_lib::Result<std::option::Option<RawTransactionIngestHttpDiscoveryWindow>> {
|
||||
let slot_count = match http_discovery_window_slot_count(start_slot, end_slot) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let capabilities = match http_role_scan_capabilities(pool, role, expected_cluster) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let strategy = match capabilities.strategy() {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let context_config = ksp_onchain_transport_lib::SolanaContextConfig::new(std::option::Option::Some(commitment), std::option::Option::None);
|
||||
let mut current_tip = std::option::Option::None;
|
||||
let discovered = match strategy {
|
||||
RawTransactionIngestHttpDiscoveryStrategy::ClosedRange => tokio::select! {
|
||||
biased;
|
||||
_ = stop_receiver.changed() => {
|
||||
return std::result::Result::Ok(std::option::Option::None);
|
||||
}
|
||||
result = pool.get_blocks(role, start_slot, std::option::Option::Some(end_slot), std::option::Option::Some(&context_config)) => result,
|
||||
},
|
||||
RawTransactionIngestHttpDiscoveryStrategy::WithLimitBoundary => {
|
||||
let tip = tokio::select! {
|
||||
biased;
|
||||
_ = stop_receiver.changed() => {
|
||||
return std::result::Result::Ok(std::option::Option::None);
|
||||
}
|
||||
result = pool.get_slot(role, std::option::Option::Some(&context_config)) => result,
|
||||
};
|
||||
let tip = match tip {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(source_transport_error(error.code())),
|
||||
};
|
||||
current_tip = std::option::Option::Some(tip);
|
||||
tokio::select! {
|
||||
biased;
|
||||
_ = stop_receiver.changed() => {
|
||||
return std::result::Result::Ok(std::option::Option::None);
|
||||
}
|
||||
result = pool.get_blocks_with_limit(role, start_slot, slot_count, std::option::Option::Some(&context_config)) => result,
|
||||
}
|
||||
},
|
||||
};
|
||||
let discovered = match discovered {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(source_transport_error(error.code())),
|
||||
};
|
||||
let validated = validate_http_block_discovery_result(start_slot, end_slot, strategy, current_tip, discovered.as_slice());
|
||||
return match validated {
|
||||
std::result::Result::Ok(value) => std::result::Result::Ok(std::option::Option::Some(value)),
|
||||
std::result::Result::Err(error) => std::result::Result::Err(error),
|
||||
};
|
||||
}
|
||||
|
||||
fn validate_http_block_polling_discovery(next_scan_slot: u64, discovered: &[u64]) -> ksp_core_lib::Result<()> {
|
||||
let mut previous = std::option::Option::None;
|
||||
for slot in discovered {
|
||||
|
||||
Reference in New Issue
Block a user