diff --git a/Cargo.toml b/Cargo.toml index 7544bb0..6c38022 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,7 +6,7 @@ resolver = "3" members = ["crates/ksp-app-config-desk", "crates/ksp-app-wallet-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-logging-lib", "crates/ksp-onchain-transport-lib", "crates/ksp-wallet-lib"] [workspace.package] -version = "0.2.9-pre.6" +version = "0.2.9-pre.7" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-onchain-transport-lib/src/grpc_subscribe.rs b/crates/ksp-onchain-transport-lib/src/grpc_subscribe.rs index 901f733..be38a58 100644 --- a/crates/ksp-onchain-transport-lib/src/grpc_subscribe.rs +++ b/crates/ksp-onchain-transport-lib/src/grpc_subscribe.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/grpc_subscribe.rs -// version: 4 +// version: 5 #[cfg(test)] const MAX_GRPC_SUBSCRIBE_ACCOUNT_DATA_LENGTH_BYTES: usize = 512 * 1024 * 1024; @@ -14,10 +14,26 @@ const MAX_GRPC_SUBSCRIBE_MEMCMP_BYTES: usize = 1024 * 1024; const MAX_GRPC_SUBSCRIBE_MEMCMP_TEXT_LENGTH_BYTES: usize = 2 * 1024 * 1024; #[cfg(test)] const MAX_GRPC_SUBSCRIBE_SLOT_DEAD_ERROR_LENGTH_BYTES: usize = 16 * 1024; +const MAX_GRPC_SUBSCRIBE_TRANSACTION_SIGNATURE_TEXT_LENGTH_BYTES: usize = 128; #[cfg(test)] const MAX_GRPC_SUBSCRIBE_UPDATE_FILTER_COUNT: usize = 1_024; #[cfg(test)] -const YELLOWSTONE_TRANSACTION_SIGNATURE_LENGTH_BYTES: usize = 64; +const MAX_GRPC_TRANSACTION_ERROR_BYTES: usize = 64 * 1024; +#[cfg(test)] +const MAX_GRPC_TRANSACTION_INSTRUCTION_DATA_BYTES: usize = 1024 * 1024; +#[cfg(test)] +const MAX_GRPC_TRANSACTION_LOG_COUNT: usize = 16_384; +#[cfg(test)] +const MAX_GRPC_TRANSACTION_LOG_LENGTH_BYTES: usize = 64 * 1024; +#[cfg(test)] +const MAX_GRPC_TRANSACTION_RETURN_DATA_BYTES: usize = 1024 * 1024; +#[cfg(test)] +const MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES: usize = 64 * 1024; +#[cfg(test)] +const MAX_GRPC_TRANSACTION_VECTOR_COUNT: usize = 65_536; +#[cfg(test)] +const YELLOWSTONE_HASH_LENGTH_BYTES: usize = 32; +const YELLOWSTONE_TRANSACTION_SIGNATURE_WIRE_LENGTH_BYTES: usize = 64; /// Validated logical filter name used by the standard Yellowstone `SubscribeRequest` maps. /// @@ -603,7 +619,7 @@ impl YellowstoneUpdateTimestamp { } /// Fixed-width transaction signature attached to a Yellowstone account update when available. -#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +#[derive(Clone, Copy, Eq, Hash, PartialEq)] pub struct YellowstoneTransactionSignature { bytes: [u8; 64], } @@ -622,6 +638,12 @@ impl YellowstoneTransactionSignature { } } +impl std::fmt::Debug for YellowstoneTransactionSignature { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter.write_str("YellowstoneTransactionSignature()"); + } +} + /// Typed account payload carried by one standard Yellowstone account update. #[derive(Clone, Eq, PartialEq)] pub struct YellowstoneAccountInfo { @@ -836,30 +858,1096 @@ impl std::fmt::Debug for YellowstoneSlotUpdate { } } -/// Transaction-family filter-group shell shared by `transactions` and `transactions_status`. -/// -/// Transaction selectors, Cuckoo filters and token-account expansion are added by `0.2.9-pre.006`. -#[derive(Clone, Debug, Default, Eq, PartialEq)] +/// Optional token-account owner expansion applied by current Yellowstone transaction filters. +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +pub enum YellowstoneTokenAccountExpansion { + /// Match owners present in either pre- or post-token balances. + All, + /// Match owners whose token balance changed or whose token account was closed. + BalanceChanged, +} + +impl YellowstoneTokenAccountExpansion { + #[cfg(test)] + const fn to_wire(self) -> i32 { + return match self { + Self::All => yellowstone_grpc_proto::geyser::TokenAccountExpansionControlFlag::All as i32, + Self::BalanceChanged => yellowstone_grpc_proto::geyser::TokenAccountExpansionControlFlag::BalanceChanged as i32, + }; + } +} + +/// Validated base58 transaction-signature selector used by Yellowstone transaction filters. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneTransactionSignatureSelector { + value: std::string::String, +} + +impl YellowstoneTransactionSignatureSelector { + /// Validates one non-empty, bounded base58 signature selector without logging its value. + pub fn new(value: impl std::convert::Into) -> ksp_core_lib::Result { + let value = value.into(); + if value.is_empty() + || value.len() > MAX_GRPC_SUBSCRIBE_TRANSACTION_SIGNATURE_TEXT_LENGTH_BYTES + || base58_decoded_length(value.as_str()) != std::option::Option::Some(YELLOWSTONE_TRANSACTION_SIGNATURE_WIRE_LENGTH_BYTES) + { + return std::result::Result::Err( + ksp_core_lib::Error::new(crate::ERROR_CODE_INVALID_RPC_PARAMETERS, "Yellowstone transaction signature selector is invalid") + .with_context("field", "grpc_subscribe.transaction.signature"), + ); + } + return std::result::Result::Ok(Self { value }); + } + + /// Returns the validated signature text. Treat this value as caller-provided filter material. + #[must_use] + pub fn as_str(&self) -> &str { + return self.value.as_str(); + } +} + +impl std::fmt::Debug for YellowstoneTransactionSignatureSelector { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter.write_str("YellowstoneTransactionSignatureSelector()"); + } +} + +/// Complete current transaction-family filter shared by `transactions` and `transactions_status`. +#[derive(Clone, Eq, PartialEq)] pub struct YellowstoneSubscribeTransactionFilter { - _private: (), + vote: std::option::Option, + failed: std::option::Option, + signature: std::option::Option, + account_include: std::vec::Vec, + account_exclude: std::vec::Vec, + account_required: std::vec::Vec, + cuckoo_account_include: std::option::Option, + token_accounts: std::option::Option, } impl YellowstoneSubscribeTransactionFilter { /// Creates an empty transaction filter group. #[must_use] - pub const fn new() -> Self { - return Self { _private: () }; + pub fn new() -> Self { + return Self { + vote: std::option::Option::None, + failed: std::option::Option::None, + signature: std::option::Option::None, + account_include: std::vec::Vec::new(), + account_exclude: std::vec::Vec::new(), + account_required: std::vec::Vec::new(), + cuckoo_account_include: std::option::Option::None, + token_accounts: std::option::Option::None, + }; + } + + /// Sets or clears the optional vote-transaction selector. + pub fn set_vote(&mut self, value: std::option::Option) { + self.vote = value; + } + + /// Sets or clears the optional failed-transaction selector. + pub fn set_failed(&mut self, value: std::option::Option) { + self.failed = value; + } + + /// Sets or clears the exact transaction-signature selector. + pub fn set_signature(&mut self, value: std::option::Option) { + self.signature = value; + } + + /// Adds one account include selector while preserving insertion order. + pub fn push_account_include(&mut self, value: ksp_core_lib::Pubkey) -> ksp_core_lib::Result<()> { + return push_transaction_selector(&mut self.account_include, value, "grpc_subscribe.transaction.account_include"); + } + + /// Adds one account exclude selector while preserving insertion order. + pub fn push_account_exclude(&mut self, value: ksp_core_lib::Pubkey) -> ksp_core_lib::Result<()> { + return push_transaction_selector(&mut self.account_exclude, value, "grpc_subscribe.transaction.account_exclude"); + } + + /// Adds one required account selector while preserving insertion order. + pub fn push_account_required(&mut self, value: ksp_core_lib::Pubkey) -> ksp_core_lib::Result<()> { + return push_transaction_selector(&mut self.account_required, value, "grpc_subscribe.transaction.account_required"); + } + + /// Sets or clears the optional compressed account-include filter. + pub fn set_cuckoo_account_include(&mut self, value: std::option::Option) { + self.cuckoo_account_include = value; + } + + /// Sets or clears token-account owner expansion. + pub fn set_token_accounts(&mut self, value: std::option::Option) { + self.token_accounts = value; + } + + /// Returns the optional vote selector. + #[must_use] + pub const fn vote(&self) -> std::option::Option { + return self.vote; + } + + /// Returns the optional failed selector. + #[must_use] + pub const fn failed(&self) -> std::option::Option { + return self.failed; + } + + /// Returns the optional exact signature selector. + #[must_use] + pub const fn signature(&self) -> std::option::Option<&crate::YellowstoneTransactionSignatureSelector> { + return self.signature.as_ref(); + } + + /// Returns ordered account include selectors. + #[must_use] + pub fn account_include(&self) -> &[ksp_core_lib::Pubkey] { + return self.account_include.as_slice(); + } + + /// Returns ordered account exclude selectors. + #[must_use] + pub fn account_exclude(&self) -> &[ksp_core_lib::Pubkey] { + return self.account_exclude.as_slice(); + } + + /// Returns ordered required account selectors. + #[must_use] + pub fn account_required(&self) -> &[ksp_core_lib::Pubkey] { + return self.account_required.as_slice(); + } + + /// Returns the optional compressed include filter. + #[must_use] + pub const fn cuckoo_account_include(&self) -> std::option::Option<&crate::YellowstoneCuckooFilter> { + return self.cuckoo_account_include.as_ref(); + } + + /// Returns optional token-account owner expansion. + #[must_use] + pub const fn token_accounts(&self) -> std::option::Option { + return self.token_accounts; } #[cfg(test)] fn to_wire(&self) -> yellowstone_grpc_proto::geyser::SubscribeRequestFilterTransactions { - return yellowstone_grpc_proto::geyser::SubscribeRequestFilterTransactions::default(); + return yellowstone_grpc_proto::geyser::SubscribeRequestFilterTransactions { + vote: self.vote, + failed: self.failed, + signature: self.signature.as_ref().map(|value| return value.as_str().to_owned()), + account_include: self.account_include.iter().map(std::string::ToString::to_string).collect(), + account_exclude: self.account_exclude.iter().map(std::string::ToString::to_string).collect(), + account_required: self.account_required.iter().map(std::string::ToString::to_string).collect(), + cuckoo_account_include: self.cuckoo_account_include.as_ref().map(crate::YellowstoneCuckooFilter::to_wire), + token_accounts: self.token_accounts.map(crate::YellowstoneTokenAccountExpansion::to_wire), + }; + } +} + +impl std::default::Default for YellowstoneSubscribeTransactionFilter { + fn default() -> Self { + return Self::new(); + } +} + +impl std::fmt::Debug for YellowstoneSubscribeTransactionFilter { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneSubscribeTransactionFilter") + .field("vote", &self.vote) + .field("failed", &self.failed) + .field("has_signature", &self.signature.is_some()) + .field("account_include_count", &self.account_include.len()) + .field("account_exclude_count", &self.account_exclude.len()) + .field("account_required_count", &self.account_required.len()) + .field("has_cuckoo_account_include", &self.cuckoo_account_include.is_some()) + .field("token_accounts", &self.token_accounts) + .finish(); + } +} + +/// Fixed-width 32-byte hash used by the Solana storage protobuf transaction wire. +#[derive(Clone, Copy, Eq, Hash, PartialEq)] +pub struct YellowstoneHashBytes { + bytes: [u8; 32], +} + +impl YellowstoneHashBytes { + /// Returns the exact 32 hash bytes. + #[must_use] + pub const fn as_bytes(&self) -> &[u8; 32] { + return &self.bytes; + } +} + +impl std::fmt::Debug for YellowstoneHashBytes { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter.write_str("YellowstoneHashBytes()"); + } +} + +/// Opaque runtime transaction error bytes carried by `solana-storage.proto`. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneTransactionError { + bytes: std::vec::Vec, +} + +impl YellowstoneTransactionError { + /// Returns the exact opaque error bytes without interpreting runtime/program semantics. + #[must_use] + pub fn as_bytes(&self) -> &[u8] { + return self.bytes.as_slice(); + } +} + +impl std::fmt::Debug for YellowstoneTransactionError { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter.debug_struct("YellowstoneTransactionError").field("length", &self.bytes.len()).finish(); + } +} + +/// Solana transaction message header projected from the storage protobuf. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct YellowstoneTransactionMessageHeader { + num_required_signatures: u32, + num_readonly_signed_accounts: u32, + num_readonly_unsigned_accounts: u32, +} + +impl YellowstoneTransactionMessageHeader { + /// Returns the required signature count. + #[must_use] + pub const fn num_required_signatures(self) -> u32 { + return self.num_required_signatures; + } + + /// Returns the readonly signed-account count. + #[must_use] + pub const fn num_readonly_signed_accounts(self) -> u32 { + return self.num_readonly_signed_accounts; + } + + /// Returns the readonly unsigned-account count. + #[must_use] + pub const fn num_readonly_unsigned_accounts(self) -> u32 { + return self.num_readonly_unsigned_accounts; + } +} + +/// One compiled instruction from the Solana storage protobuf transaction message. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneCompiledInstruction { + program_id_index: u32, + accounts: std::vec::Vec, + data: std::vec::Vec, +} + +impl YellowstoneCompiledInstruction { + /// Returns the program-id index. + #[must_use] + pub const fn program_id_index(&self) -> u32 { + return self.program_id_index; + } + + /// Returns the exact account-index bytes. + #[must_use] + pub fn accounts(&self) -> &[u8] { + return self.accounts.as_slice(); + } + + /// Returns the exact instruction data bytes. + #[must_use] + pub fn data(&self) -> &[u8] { + return self.data.as_slice(); + } +} + +impl std::fmt::Debug for YellowstoneCompiledInstruction { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneCompiledInstruction") + .field("program_id_index", &self.program_id_index) + .field("account_index_count", &self.accounts.len()) + .field("data_length", &self.data.len()) + .finish(); + } +} + +/// One address-table lookup from a versioned Solana transaction message. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneMessageAddressTableLookup { + account_key: ksp_core_lib::Pubkey, + writable_indexes: std::vec::Vec, + readonly_indexes: std::vec::Vec, +} + +impl YellowstoneMessageAddressTableLookup { + /// Returns the lookup-table account key. + #[must_use] + pub const fn account_key(&self) -> &ksp_core_lib::Pubkey { + return &self.account_key; + } + + /// Returns ordered writable lookup indexes. + #[must_use] + pub fn writable_indexes(&self) -> &[u8] { + return self.writable_indexes.as_slice(); + } + + /// Returns ordered readonly lookup indexes. + #[must_use] + pub fn readonly_indexes(&self) -> &[u8] { + return self.readonly_indexes.as_slice(); + } +} + +impl std::fmt::Debug for YellowstoneMessageAddressTableLookup { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneMessageAddressTableLookup") + .field("writable_index_count", &self.writable_indexes.len()) + .field("readonly_index_count", &self.readonly_indexes.len()) + .finish_non_exhaustive(); + } +} + +/// Optional inline budget configuration introduced by Solana Transaction V1 / SIMD-0385. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct YellowstoneTransactionConfig { + priority_fee: std::option::Option, + compute_unit_limit: std::option::Option, + loaded_accounts_data_size_limit: std::option::Option, + heap_size: std::option::Option, +} + +impl YellowstoneTransactionConfig { + /// Returns the optional priority fee. + #[must_use] + pub const fn priority_fee(self) -> std::option::Option { + return self.priority_fee; + } + + /// Returns the optional compute-unit limit. + #[must_use] + pub const fn compute_unit_limit(self) -> std::option::Option { + return self.compute_unit_limit; + } + + /// Returns the optional loaded-account-data-size limit. + #[must_use] + pub const fn loaded_accounts_data_size_limit(self) -> std::option::Option { + return self.loaded_accounts_data_size_limit; + } + + /// Returns the optional heap-size override. + #[must_use] + pub const fn heap_size(self) -> std::option::Option { + return self.heap_size; + } +} + +/// Complete Solana transaction message projected from current `solana-storage.proto`. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneTransactionMessage { + header: crate::YellowstoneTransactionMessageHeader, + account_keys: std::vec::Vec, + recent_blockhash: crate::YellowstoneHashBytes, + instructions: std::vec::Vec, + versioned: bool, + address_table_lookups: std::vec::Vec, + config: std::option::Option, +} + +impl YellowstoneTransactionMessage { + /// Returns the message header. + #[must_use] + pub const fn header(&self) -> crate::YellowstoneTransactionMessageHeader { + return self.header; + } + + /// Returns ordered static account keys. + #[must_use] + pub fn account_keys(&self) -> &[ksp_core_lib::Pubkey] { + return self.account_keys.as_slice(); + } + + /// Returns the recent blockhash bytes. + #[must_use] + pub const fn recent_blockhash(&self) -> crate::YellowstoneHashBytes { + return self.recent_blockhash; + } + + /// Returns ordered compiled instructions. + #[must_use] + pub fn instructions(&self) -> &[crate::YellowstoneCompiledInstruction] { + return self.instructions.as_slice(); + } + + /// Returns the protobuf `versioned` marker. + #[must_use] + pub const fn versioned(&self) -> bool { + return self.versioned; + } + + /// Returns ordered address-table lookups. + #[must_use] + pub fn address_table_lookups(&self) -> &[crate::YellowstoneMessageAddressTableLookup] { + return self.address_table_lookups.as_slice(); + } + + /// Returns optional Transaction V1 inline budget configuration. + #[must_use] + pub const fn config(&self) -> std::option::Option { + return self.config; + } +} + +impl std::fmt::Debug for YellowstoneTransactionMessage { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneTransactionMessage") + .field("header", &self.header) + .field("account_key_count", &self.account_keys.len()) + .field("instruction_count", &self.instructions.len()) + .field("versioned", &self.versioned) + .field("address_table_lookup_count", &self.address_table_lookups.len()) + .field("config", &self.config) + .finish(); + } +} + +/// Solana transaction body carried by Yellowstone storage protobuf messages. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneStoredTransaction { + signatures: std::vec::Vec, + message: crate::YellowstoneTransactionMessage, +} + +impl YellowstoneStoredTransaction { + /// Returns ordered transaction signatures. + #[must_use] + pub fn signatures(&self) -> &[crate::YellowstoneTransactionSignature] { + return self.signatures.as_slice(); + } + + /// Returns the decoded transaction message. + #[must_use] + pub const fn message(&self) -> &crate::YellowstoneTransactionMessage { + return &self.message; + } +} + +impl std::fmt::Debug for YellowstoneStoredTransaction { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneStoredTransaction") + .field("signature_count", &self.signatures.len()) + .field("message", &self.message) + .finish(); + } +} + +/// One inner instruction from transaction status metadata. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneInnerInstruction { + program_id_index: u32, + accounts: std::vec::Vec, + data: std::vec::Vec, + stack_height: std::option::Option, +} + +impl YellowstoneInnerInstruction { + /// Returns the program-id index. + #[must_use] + pub const fn program_id_index(&self) -> u32 { + return self.program_id_index; + } + + /// Returns account-index bytes. + #[must_use] + pub fn accounts(&self) -> &[u8] { + return self.accounts.as_slice(); + } + + /// Returns instruction data bytes. + #[must_use] + pub fn data(&self) -> &[u8] { + return self.data.as_slice(); + } + + /// Returns optional invocation stack height. + #[must_use] + pub const fn stack_height(&self) -> std::option::Option { + return self.stack_height; + } +} + +impl std::fmt::Debug for YellowstoneInnerInstruction { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneInnerInstruction") + .field("program_id_index", &self.program_id_index) + .field("account_index_count", &self.accounts.len()) + .field("data_length", &self.data.len()) + .field("stack_height", &self.stack_height) + .finish(); + } +} + +/// One indexed inner-instruction group from transaction status metadata. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct YellowstoneInnerInstructions { + index: u32, + instructions: std::vec::Vec, +} + +impl YellowstoneInnerInstructions { + /// Returns the outer instruction index. + #[must_use] + pub const fn index(&self) -> u32 { + return self.index; + } + + /// Returns ordered inner instructions. + #[must_use] + pub fn instructions(&self) -> &[crate::YellowstoneInnerInstruction] { + return self.instructions.as_slice(); + } +} + +/// UI token amount from current Solana storage protobuf metadata. +#[derive(Clone, PartialEq)] +pub struct YellowstoneUiTokenAmount { + ui_amount: f64, + decimals: u32, + amount: std::string::String, + ui_amount_string: std::string::String, +} + +impl YellowstoneUiTokenAmount { + /// Returns the floating UI amount carried by the protobuf. + #[must_use] + pub const fn ui_amount(&self) -> f64 { + return self.ui_amount; + } + + /// Returns token decimals. + #[must_use] + pub const fn decimals(&self) -> u32 { + return self.decimals; + } + + /// Returns the exact integer amount string. + #[must_use] + pub fn amount(&self) -> &str { + return self.amount.as_str(); + } + + /// Returns the exact UI amount string. + #[must_use] + pub fn ui_amount_string(&self) -> &str { + return self.ui_amount_string.as_str(); + } +} + +impl std::fmt::Debug for YellowstoneUiTokenAmount { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneUiTokenAmount") + .field("decimals", &self.decimals) + .field("has_amount", &!self.amount.is_empty()) + .field("has_ui_amount_string", &!self.ui_amount_string.is_empty()) + .finish_non_exhaustive(); + } +} + +/// One pre/post token balance entry from transaction status metadata. +#[derive(Clone, PartialEq)] +pub struct YellowstoneTokenBalance { + account_index: u32, + mint: std::string::String, + ui_token_amount: std::option::Option, + owner: std::string::String, + program_id: std::string::String, +} + +impl YellowstoneTokenBalance { + /// Returns the account index. + #[must_use] + pub const fn account_index(&self) -> u32 { + return self.account_index; + } + + /// Returns the exact mint text. + #[must_use] + pub fn mint(&self) -> &str { + return self.mint.as_str(); + } + + /// Returns optional UI token amount data exactly as represented by protobuf message presence. + #[must_use] + pub const fn ui_token_amount(&self) -> std::option::Option<&crate::YellowstoneUiTokenAmount> { + return self.ui_token_amount.as_ref(); + } + + /// Returns the exact owner text, including an empty legacy value when present on the wire. + #[must_use] + pub fn owner(&self) -> &str { + return self.owner.as_str(); + } + + /// Returns the exact token-program id text, including an empty legacy value. + #[must_use] + pub fn program_id(&self) -> &str { + return self.program_id.as_str(); + } +} + +impl std::fmt::Debug for YellowstoneTokenBalance { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneTokenBalance") + .field("account_index", &self.account_index) + .field("has_ui_token_amount", &self.ui_token_amount.is_some()) + .field("has_owner", &!self.owner.is_empty()) + .field("has_program_id", &!self.program_id.is_empty()) + .finish_non_exhaustive(); + } +} + +/// Return-data payload from transaction status metadata. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneReturnData { + program_id: ksp_core_lib::Pubkey, + data: std::vec::Vec, +} + +impl YellowstoneReturnData { + /// Returns the program id that produced the return data. + #[must_use] + pub const fn program_id(&self) -> &ksp_core_lib::Pubkey { + return &self.program_id; + } + + /// Returns the exact return bytes. + #[must_use] + pub fn data(&self) -> &[u8] { + return self.data.as_slice(); + } +} + +impl std::fmt::Debug for YellowstoneReturnData { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter.debug_struct("YellowstoneReturnData").field("data_length", &self.data.len()).finish_non_exhaustive(); + } +} + +/// Reward classification carried by current Solana storage protobuf metadata. +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +pub enum YellowstoneRewardType { + /// Unspecified reward type. + Unspecified, + /// Fee reward. + Fee, + /// Rent reward. + Rent, + /// Staking reward. + Staking, + /// Voting reward. + Voting, + /// Deactivated stake reward. + DeactivatedStake, + /// Explicitly preserved future/unknown numeric value. + Unknown(i32), +} + +/// One reward entry from transaction status metadata. +#[derive(Clone, Eq, PartialEq)] +pub struct YellowstoneReward { + pubkey: ksp_core_lib::Pubkey, + lamports: i64, + post_balance: u64, + reward_type: crate::YellowstoneRewardType, + commission: std::string::String, + commission_bps: std::string::String, +} + +impl YellowstoneReward { + /// Returns the reward account. + #[must_use] + pub const fn pubkey(&self) -> &ksp_core_lib::Pubkey { + return &self.pubkey; + } + + /// Returns signed reward lamports. + #[must_use] + pub const fn lamports(&self) -> i64 { + return self.lamports; + } + + /// Returns post-reward balance. + #[must_use] + pub const fn post_balance(&self) -> u64 { + return self.post_balance; + } + + /// Returns reward type, preserving unknown numeric values explicitly. + #[must_use] + pub const fn reward_type(&self) -> crate::YellowstoneRewardType { + return self.reward_type; + } + + /// Returns legacy commission text exactly as carried on the protobuf wire. + #[must_use] + pub fn commission(&self) -> &str { + return self.commission.as_str(); + } + + /// Returns commission basis-point text exactly as carried on the protobuf wire. + #[must_use] + pub fn commission_bps(&self) -> &str { + return self.commission_bps.as_str(); + } +} + +impl std::fmt::Debug for YellowstoneReward { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneReward") + .field("lamports", &self.lamports) + .field("post_balance", &self.post_balance) + .field("reward_type", &self.reward_type) + .field("has_commission", &!self.commission.is_empty()) + .field("has_commission_bps", &!self.commission_bps.is_empty()) + .finish_non_exhaustive(); + } +} + +/// Complete transaction status metadata projected from current `solana-storage.proto`. +#[derive(Clone, PartialEq)] +pub struct YellowstoneTransactionStatusMeta { + error: std::option::Option, + fee: u64, + pre_balances: std::vec::Vec, + post_balances: std::vec::Vec, + inner_instructions: std::vec::Vec, + inner_instructions_none: bool, + log_messages: std::vec::Vec, + log_messages_none: bool, + pre_token_balances: std::vec::Vec, + post_token_balances: std::vec::Vec, + rewards: std::vec::Vec, + loaded_writable_addresses: std::vec::Vec, + loaded_readonly_addresses: std::vec::Vec, + return_data: std::option::Option, + return_data_none: bool, + compute_units_consumed: std::option::Option, + cost_units: std::option::Option, +} + +impl YellowstoneTransactionStatusMeta { + /// Returns optional opaque runtime error bytes. + #[must_use] + pub const fn error(&self) -> std::option::Option<&crate::YellowstoneTransactionError> { + return self.error.as_ref(); + } + + /// Returns transaction fee in lamports. + #[must_use] + pub const fn fee(&self) -> u64 { + return self.fee; + } + + /// Returns pre-transaction lamport balances. + #[must_use] + pub fn pre_balances(&self) -> &[u64] { + return self.pre_balances.as_slice(); + } + + /// Returns post-transaction lamport balances. + #[must_use] + pub fn post_balances(&self) -> &[u64] { + return self.post_balances.as_slice(); + } + + /// Returns inner-instruction groups. + #[must_use] + pub fn inner_instructions(&self) -> &[crate::YellowstoneInnerInstructions] { + return self.inner_instructions.as_slice(); + } + + /// Returns the explicit legacy `inner_instructions_none` marker. + #[must_use] + pub const fn inner_instructions_none(&self) -> bool { + return self.inner_instructions_none; + } + + /// Returns ordered runtime log messages. + #[must_use] + pub fn log_messages(&self) -> &[std::string::String] { + return self.log_messages.as_slice(); + } + + /// Returns the explicit legacy `log_messages_none` marker. + #[must_use] + pub const fn log_messages_none(&self) -> bool { + return self.log_messages_none; + } + + /// Returns pre-token balances. + #[must_use] + pub fn pre_token_balances(&self) -> &[crate::YellowstoneTokenBalance] { + return self.pre_token_balances.as_slice(); + } + + /// Returns post-token balances. + #[must_use] + pub fn post_token_balances(&self) -> &[crate::YellowstoneTokenBalance] { + return self.post_token_balances.as_slice(); + } + + /// Returns reward entries. + #[must_use] + pub fn rewards(&self) -> &[crate::YellowstoneReward] { + return self.rewards.as_slice(); + } + + /// Returns loaded writable addresses. + #[must_use] + pub fn loaded_writable_addresses(&self) -> &[ksp_core_lib::Pubkey] { + return self.loaded_writable_addresses.as_slice(); + } + + /// Returns loaded readonly addresses. + #[must_use] + pub fn loaded_readonly_addresses(&self) -> &[ksp_core_lib::Pubkey] { + return self.loaded_readonly_addresses.as_slice(); + } + + /// Returns optional return data while preserving the separate protobuf `return_data_none` marker. + #[must_use] + pub const fn return_data(&self) -> std::option::Option<&crate::YellowstoneReturnData> { + return self.return_data.as_ref(); + } + + /// Returns the explicit legacy `return_data_none` marker. + #[must_use] + pub const fn return_data_none(&self) -> bool { + return self.return_data_none; + } + + /// Returns optional compute units consumed. + #[must_use] + pub const fn compute_units_consumed(&self) -> std::option::Option { + return self.compute_units_consumed; + } + + /// Returns optional total transaction cost units. + #[must_use] + pub const fn cost_units(&self) -> std::option::Option { + return self.cost_units; + } +} + +impl std::fmt::Debug for YellowstoneTransactionStatusMeta { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneTransactionStatusMeta") + .field("has_error", &self.error.is_some()) + .field("fee", &self.fee) + .field("pre_balance_count", &self.pre_balances.len()) + .field("post_balance_count", &self.post_balances.len()) + .field("inner_instruction_group_count", &self.inner_instructions.len()) + .field("inner_instructions_none", &self.inner_instructions_none) + .field("log_message_count", &self.log_messages.len()) + .field("log_messages_none", &self.log_messages_none) + .field("pre_token_balance_count", &self.pre_token_balances.len()) + .field("post_token_balance_count", &self.post_token_balances.len()) + .field("reward_count", &self.rewards.len()) + .field("loaded_writable_address_count", &self.loaded_writable_addresses.len()) + .field("loaded_readonly_address_count", &self.loaded_readonly_addresses.len()) + .field("has_return_data", &self.return_data.is_some()) + .field("return_data_none", &self.return_data_none) + .field("compute_units_consumed", &self.compute_units_consumed) + .field("cost_units", &self.cost_units) + .finish(); + } +} + +/// Full transaction info carried by `SubscribeUpdateTransaction` and later reused by block updates. +#[derive(Clone, PartialEq)] +pub struct YellowstoneTransactionInfo { + signature: crate::YellowstoneTransactionSignature, + is_vote: bool, + transaction: crate::YellowstoneStoredTransaction, + meta: crate::YellowstoneTransactionStatusMeta, + index: u64, +} + +impl YellowstoneTransactionInfo { + /// Returns the primary transaction signature. + #[must_use] + pub const fn signature(&self) -> crate::YellowstoneTransactionSignature { + return self.signature; + } + + /// Returns whether Yellowstone classifies this as a vote transaction. + #[must_use] + pub const fn is_vote(&self) -> bool { + return self.is_vote; + } + + /// Returns the complete typed transaction body. + #[must_use] + pub const fn transaction(&self) -> &crate::YellowstoneStoredTransaction { + return &self.transaction; + } + + /// Returns complete typed transaction status metadata. + #[must_use] + pub const fn meta(&self) -> &crate::YellowstoneTransactionStatusMeta { + return &self.meta; + } + + /// Returns the transaction index within the block. + #[must_use] + pub const fn index(&self) -> u64 { + return self.index; + } +} + +impl std::fmt::Debug for YellowstoneTransactionInfo { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneTransactionInfo") + .field("is_vote", &self.is_vote) + .field("transaction", &self.transaction) + .field("meta", &self.meta) + .field("index", &self.index) + .finish_non_exhaustive(); + } +} + +/// Full standard Yellowstone transaction update projected into KSP-owned types. +#[derive(Clone, PartialEq)] +pub struct YellowstoneTransactionUpdate { + filters: std::vec::Vec, + created_at: std::option::Option, + transaction: crate::YellowstoneTransactionInfo, + slot: u64, +} + +impl YellowstoneTransactionUpdate { + /// Returns echoed matching filter names in server order. + #[must_use] + pub fn filters(&self) -> &[crate::YellowstoneSubscribeFilterName] { + return self.filters.as_slice(); + } + + /// Returns optional server creation timestamp. + #[must_use] + pub const fn created_at(&self) -> std::option::Option { + return self.created_at; + } + + /// Returns complete transaction info. + #[must_use] + pub const fn transaction(&self) -> &crate::YellowstoneTransactionInfo { + return &self.transaction; + } + + /// Returns the slot containing the transaction. + #[must_use] + pub const fn slot(&self) -> u64 { + return self.slot; + } +} + +impl std::fmt::Debug for YellowstoneTransactionUpdate { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneTransactionUpdate") + .field("filter_count", &self.filters.len()) + .field("created_at", &self.created_at) + .field("transaction", &self.transaction) + .field("slot", &self.slot) + .finish(); + } +} + +/// Lightweight standard Yellowstone transaction-status update. +#[derive(Clone, PartialEq)] +pub struct YellowstoneTransactionStatusUpdate { + filters: std::vec::Vec, + created_at: std::option::Option, + slot: u64, + signature: crate::YellowstoneTransactionSignature, + is_vote: bool, + index: u64, + error: std::option::Option, +} + +impl YellowstoneTransactionStatusUpdate { + /// Returns echoed matching filter names in server order. + #[must_use] + pub fn filters(&self) -> &[crate::YellowstoneSubscribeFilterName] { + return self.filters.as_slice(); + } + + /// Returns optional server creation timestamp. + #[must_use] + pub const fn created_at(&self) -> std::option::Option { + return self.created_at; + } + + /// Returns the containing slot. + #[must_use] + pub const fn slot(&self) -> u64 { + return self.slot; + } + + /// Returns the transaction signature. + #[must_use] + pub const fn signature(&self) -> crate::YellowstoneTransactionSignature { + return self.signature; + } + + /// Returns whether this is a vote transaction. + #[must_use] + pub const fn is_vote(&self) -> bool { + return self.is_vote; + } + + /// Returns the transaction index within the block. + #[must_use] + pub const fn index(&self) -> u64 { + return self.index; + } + + /// Returns optional opaque runtime error bytes. + #[must_use] + pub const fn error(&self) -> std::option::Option<&crate::YellowstoneTransactionError> { + return self.error.as_ref(); + } +} + +impl std::fmt::Debug for YellowstoneTransactionStatusUpdate { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter + .debug_struct("YellowstoneTransactionStatusUpdate") + .field("filter_count", &self.filters.len()) + .field("created_at", &self.created_at) + .field("slot", &self.slot) + .field("is_vote", &self.is_vote) + .field("index", &self.index) + .field("has_error", &self.error.is_some()) + .finish_non_exhaustive(); } } /// Block-family filter-group shell for the standard Yellowstone subscribe surface. /// -/// Block selectors, include flags and Cuckoo filters are added by `0.2.9-pre.007`. +/// Block selectors, include flags and Cuckoo filters are added by `0.2.9-pre.008`. #[derive(Clone, Debug, Default, Eq, PartialEq)] pub struct YellowstoneSubscribeBlockFilter { _private: (), @@ -920,7 +2008,8 @@ impl YellowstoneSubscribeEntryFilter { /// /// The seven upstream maps are represented independently and retain named empty entries. An entirely empty map is the logical KSP representation of no active /// filter in that family; protobuf map encoding does not distinguish an omitted map from an empty map. Filter-group names are globally unique across all seven -/// maps so the names echoed by `SubscribeUpdate.filters` remain unambiguous. `Debug` exposes only counts and common scalar options, never filter names or future +/// maps so the names echoed by `SubscribeUpdate.filters` remain unambiguous. `Debug` exposes only counts and common scalar options, never filter names or +/// future /// filter payloads. #[derive(Clone, Eq, PartialEq)] pub struct YellowstoneSubscribeRequest { @@ -938,7 +2027,8 @@ pub struct YellowstoneSubscribeRequest { } impl YellowstoneSubscribeRequest { - /// Creates an empty subscribe request. Empty requests are valid because later bidi lifecycle code uses request mutations to clear filters or carry ping state. + /// Creates an empty subscribe request. Empty requests are valid because later bidi lifecycle code uses request mutations to clear filters or carry ping + /// state. #[must_use] pub fn new() -> Self { return Self { @@ -1275,6 +2365,436 @@ fn decode_slot_update(wire: yellowstone_grpc_proto::geyser::SubscribeUpdate) -> }); } +#[cfg(test)] +fn decode_transaction_update(wire: yellowstone_grpc_proto::geyser::SubscribeUpdate) -> ksp_core_lib::Result { + let (filters, created_at) = match decode_update_envelope(wire.filters, wire.created_at) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let update = match wire.update_oneof { + std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::Transaction(value)) => value, + _ => return invalid_subscribe_response("transaction", "Yellowstone update does not contain a transaction payload"), + }; + let transaction = match update.transaction { + std::option::Option::Some(value) => match decode_transaction_info(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::option::Option::None => return invalid_subscribe_response("transaction", "Yellowstone transaction update is missing transaction info"), + }; + return std::result::Result::Ok(crate::YellowstoneTransactionUpdate { filters, created_at, transaction, slot: update.slot }); +} + +#[cfg(test)] +fn decode_transaction_status_update(wire: yellowstone_grpc_proto::geyser::SubscribeUpdate) -> ksp_core_lib::Result { + let (filters, created_at) = match decode_update_envelope(wire.filters, wire.created_at) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let update = match wire.update_oneof { + std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::TransactionStatus(value)) => value, + _ => return invalid_subscribe_response("transaction_status", "Yellowstone update does not contain a transaction-status payload"), + }; + let signature = match decode_transaction_signature("transaction_status.signature", update.signature) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let error = match update.err { + std::option::Option::Some(value) => match decode_transaction_error("transaction_status.err", value) { + std::result::Result::Ok(value) => std::option::Option::Some(value), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::option::Option::None => std::option::Option::None, + }; + return std::result::Result::Ok(crate::YellowstoneTransactionStatusUpdate { + filters, + created_at, + slot: update.slot, + signature, + is_vote: update.is_vote, + index: update.index, + error, + }); +} + +#[cfg(test)] +fn decode_transaction_info(wire: yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionInfo) -> ksp_core_lib::Result { + let signature = match decode_transaction_signature("transaction.signature", wire.signature) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let transaction = match wire.transaction { + std::option::Option::Some(value) => match decode_stored_transaction(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::option::Option::None => return invalid_subscribe_response("transaction.transaction", "Yellowstone transaction info is missing transaction body"), + }; + let meta = match wire.meta { + std::option::Option::Some(value) => match decode_transaction_status_meta(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::option::Option::None => return invalid_subscribe_response("transaction.meta", "Yellowstone transaction info is missing status metadata"), + }; + return std::result::Result::Ok(crate::YellowstoneTransactionInfo { signature, is_vote: wire.is_vote, transaction, meta, index: wire.index }); +} + +#[cfg(test)] +fn decode_stored_transaction( + wire: yellowstone_grpc_proto::solana::storage::confirmed_block::Transaction, +) -> ksp_core_lib::Result { + if wire.signatures.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT { + return invalid_subscribe_response("transaction.signatures", "Yellowstone transaction signature count exceeds the KSP bound"); + } + let mut signatures = std::vec::Vec::with_capacity(wire.signatures.len()); + for signature in wire.signatures { + let signature = match decode_transaction_signature("transaction.signatures", signature) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + signatures.push(signature); + } + let message = match wire.message { + std::option::Option::Some(value) => match decode_transaction_message(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::option::Option::None => return invalid_subscribe_response("transaction.message", "Yellowstone transaction body is missing message"), + }; + return std::result::Result::Ok(crate::YellowstoneStoredTransaction { signatures, message }); +} + +#[cfg(test)] +fn decode_transaction_message( + wire: yellowstone_grpc_proto::solana::storage::confirmed_block::Message, +) -> ksp_core_lib::Result { + if wire.account_keys.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.instructions.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.address_table_lookups.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + { + return invalid_subscribe_response("transaction.message", "Yellowstone transaction message collection exceeds the KSP bound"); + } + let header = match wire.header { + std::option::Option::Some(value) => crate::YellowstoneTransactionMessageHeader { + num_required_signatures: value.num_required_signatures, + num_readonly_signed_accounts: value.num_readonly_signed_accounts, + num_readonly_unsigned_accounts: value.num_readonly_unsigned_accounts, + }, + std::option::Option::None => return invalid_subscribe_response("transaction.message.header", "Yellowstone transaction message is missing header"), + }; + let mut account_keys = std::vec::Vec::with_capacity(wire.account_keys.len()); + for value in wire.account_keys { + let value = match decode_pubkey_bytes("transaction.message.account_keys", value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + account_keys.push(value); + } + let recent_blockhash = match decode_hash_bytes("transaction.message.recent_blockhash", wire.recent_blockhash) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let mut instructions = std::vec::Vec::with_capacity(wire.instructions.len()); + for value in wire.instructions { + let value = match decode_compiled_instruction(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + instructions.push(value); + } + let mut address_table_lookups = std::vec::Vec::with_capacity(wire.address_table_lookups.len()); + for value in wire.address_table_lookups { + let value = match decode_message_address_table_lookup(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + address_table_lookups.push(value); + } + let config = wire.config.map(|value| { + return crate::YellowstoneTransactionConfig { + priority_fee: value.priority_fee, + compute_unit_limit: value.compute_unit_limit, + loaded_accounts_data_size_limit: value.loaded_accounts_data_size_limit, + heap_size: value.heap_size, + }; + }); + return std::result::Result::Ok(crate::YellowstoneTransactionMessage { + header, + account_keys, + recent_blockhash, + instructions, + versioned: wire.versioned, + address_table_lookups, + config, + }); +} + +#[cfg(test)] +fn decode_compiled_instruction( + wire: yellowstone_grpc_proto::solana::storage::confirmed_block::CompiledInstruction, +) -> ksp_core_lib::Result { + if wire.accounts.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT || wire.data.len() > MAX_GRPC_TRANSACTION_INSTRUCTION_DATA_BYTES { + return invalid_subscribe_response("transaction.message.instructions", "Yellowstone compiled instruction exceeds the KSP bound"); + } + return std::result::Result::Ok(crate::YellowstoneCompiledInstruction { + program_id_index: wire.program_id_index, + accounts: wire.accounts, + data: wire.data, + }); +} + +#[cfg(test)] +fn decode_message_address_table_lookup( + wire: yellowstone_grpc_proto::solana::storage::confirmed_block::MessageAddressTableLookup, +) -> ksp_core_lib::Result { + if wire.writable_indexes.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT || wire.readonly_indexes.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT { + return invalid_subscribe_response("transaction.message.address_table_lookups", "Yellowstone address-table index collection exceeds the KSP bound"); + } + let account_key = match decode_pubkey_bytes("transaction.message.address_table_lookup.account_key", wire.account_key) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + return std::result::Result::Ok(crate::YellowstoneMessageAddressTableLookup { + account_key, + writable_indexes: wire.writable_indexes, + readonly_indexes: wire.readonly_indexes, + }); +} + +#[cfg(test)] +fn decode_transaction_status_meta( + wire: yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionStatusMeta, +) -> ksp_core_lib::Result { + if wire.pre_balances.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.post_balances.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.inner_instructions.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.log_messages.len() > MAX_GRPC_TRANSACTION_LOG_COUNT + || wire.pre_token_balances.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.post_token_balances.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.rewards.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.loaded_writable_addresses.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + || wire.loaded_readonly_addresses.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT + { + return invalid_subscribe_response("transaction.meta", "Yellowstone transaction metadata collection exceeds the KSP bound"); + } + let error = match wire.err { + std::option::Option::Some(value) => match decode_transaction_error("transaction.meta.err", value) { + std::result::Result::Ok(value) => std::option::Option::Some(value), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::option::Option::None => std::option::Option::None, + }; + let mut inner_instructions = std::vec::Vec::with_capacity(wire.inner_instructions.len()); + for value in wire.inner_instructions { + let value = match decode_inner_instructions(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + inner_instructions.push(value); + } + for value in &wire.log_messages { + if value.len() > MAX_GRPC_TRANSACTION_LOG_LENGTH_BYTES { + return invalid_subscribe_response("transaction.meta.log_messages", "Yellowstone transaction log message exceeds the KSP bound"); + } + } + let mut pre_token_balances = std::vec::Vec::with_capacity(wire.pre_token_balances.len()); + for value in wire.pre_token_balances { + let value = match decode_token_balance(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + pre_token_balances.push(value); + } + let mut post_token_balances = std::vec::Vec::with_capacity(wire.post_token_balances.len()); + for value in wire.post_token_balances { + let value = match decode_token_balance(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + post_token_balances.push(value); + } + let mut rewards = std::vec::Vec::with_capacity(wire.rewards.len()); + for value in wire.rewards { + let value = match decode_reward(value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + rewards.push(value); + } + let mut loaded_writable_addresses = std::vec::Vec::with_capacity(wire.loaded_writable_addresses.len()); + for value in wire.loaded_writable_addresses { + let value = match decode_pubkey_bytes("transaction.meta.loaded_writable_addresses", value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + loaded_writable_addresses.push(value); + } + let mut loaded_readonly_addresses = std::vec::Vec::with_capacity(wire.loaded_readonly_addresses.len()); + for value in wire.loaded_readonly_addresses { + let value = match decode_pubkey_bytes("transaction.meta.loaded_readonly_addresses", value) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + loaded_readonly_addresses.push(value); + } + let return_data = match wire.return_data { + std::option::Option::Some(value) => match decode_return_data(value) { + std::result::Result::Ok(value) => std::option::Option::Some(value), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::option::Option::None => std::option::Option::None, + }; + return std::result::Result::Ok(crate::YellowstoneTransactionStatusMeta { + error, + fee: wire.fee, + pre_balances: wire.pre_balances, + post_balances: wire.post_balances, + inner_instructions, + inner_instructions_none: wire.inner_instructions_none, + log_messages: wire.log_messages, + log_messages_none: wire.log_messages_none, + pre_token_balances, + post_token_balances, + rewards, + loaded_writable_addresses, + loaded_readonly_addresses, + return_data, + return_data_none: wire.return_data_none, + compute_units_consumed: wire.compute_units_consumed, + cost_units: wire.cost_units, + }); +} + +#[cfg(test)] +fn decode_inner_instructions( + wire: yellowstone_grpc_proto::solana::storage::confirmed_block::InnerInstructions, +) -> ksp_core_lib::Result { + if wire.instructions.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT { + return invalid_subscribe_response("transaction.meta.inner_instructions", "Yellowstone inner-instruction group exceeds the KSP bound"); + } + let mut instructions = std::vec::Vec::with_capacity(wire.instructions.len()); + for value in wire.instructions { + if value.accounts.len() > MAX_GRPC_TRANSACTION_VECTOR_COUNT || value.data.len() > MAX_GRPC_TRANSACTION_INSTRUCTION_DATA_BYTES { + return invalid_subscribe_response("transaction.meta.inner_instructions", "Yellowstone inner instruction exceeds the KSP bound"); + } + instructions.push(crate::YellowstoneInnerInstruction { + program_id_index: value.program_id_index, + accounts: value.accounts, + data: value.data, + stack_height: value.stack_height, + }); + } + return std::result::Result::Ok(crate::YellowstoneInnerInstructions { index: wire.index, instructions }); +} + +#[cfg(test)] +fn decode_token_balance(wire: yellowstone_grpc_proto::solana::storage::confirmed_block::TokenBalance) -> ksp_core_lib::Result { + if wire.mint.len() > MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES + || wire.owner.len() > MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES + || wire.program_id.len() > MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES + { + return invalid_subscribe_response("transaction.meta.token_balances", "Yellowstone token-balance text exceeds the KSP bound"); + } + let ui_token_amount = match wire.ui_token_amount { + std::option::Option::Some(value) => { + if value.amount.len() > MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES || value.ui_amount_string.len() > MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES { + return invalid_subscribe_response("transaction.meta.token_balances", "Yellowstone token amount text exceeds the KSP bound"); + } + std::option::Option::Some(crate::YellowstoneUiTokenAmount { + ui_amount: value.ui_amount, + decimals: value.decimals, + amount: value.amount, + ui_amount_string: value.ui_amount_string, + }) + }, + std::option::Option::None => std::option::Option::None, + }; + return std::result::Result::Ok(crate::YellowstoneTokenBalance { + account_index: wire.account_index, + mint: wire.mint, + ui_token_amount, + owner: wire.owner, + program_id: wire.program_id, + }); +} + +#[cfg(test)] +fn decode_return_data(wire: yellowstone_grpc_proto::solana::storage::confirmed_block::ReturnData) -> ksp_core_lib::Result { + if wire.data.len() > MAX_GRPC_TRANSACTION_RETURN_DATA_BYTES { + return invalid_subscribe_response("transaction.meta.return_data", "Yellowstone transaction return data exceeds the KSP bound"); + } + let program_id = match decode_pubkey_bytes("transaction.meta.return_data.program_id", wire.program_id) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + return std::result::Result::Ok(crate::YellowstoneReturnData { program_id, data: wire.data }); +} + +#[cfg(test)] +fn decode_reward(wire: yellowstone_grpc_proto::solana::storage::confirmed_block::Reward) -> ksp_core_lib::Result { + if wire.pubkey.len() > MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES + || wire.commission.len() > MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES + || wire.commission_bps.len() > MAX_GRPC_TRANSACTION_TEXT_LENGTH_BYTES + { + return invalid_subscribe_response("transaction.meta.rewards", "Yellowstone reward text exceeds the KSP bound"); + } + let pubkey = match wire.pubkey.parse::() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => { + return invalid_subscribe_response("transaction.meta.rewards.pubkey", "Yellowstone reward contains an invalid public key"); + }, + }; + let reward_type = match wire.reward_type { + value if value == yellowstone_grpc_proto::solana::storage::confirmed_block::RewardType::Unspecified as i32 => crate::YellowstoneRewardType::Unspecified, + value if value == yellowstone_grpc_proto::solana::storage::confirmed_block::RewardType::Fee as i32 => crate::YellowstoneRewardType::Fee, + value if value == yellowstone_grpc_proto::solana::storage::confirmed_block::RewardType::Rent as i32 => crate::YellowstoneRewardType::Rent, + value if value == yellowstone_grpc_proto::solana::storage::confirmed_block::RewardType::Staking as i32 => crate::YellowstoneRewardType::Staking, + value if value == yellowstone_grpc_proto::solana::storage::confirmed_block::RewardType::Voting as i32 => crate::YellowstoneRewardType::Voting, + value if value == yellowstone_grpc_proto::solana::storage::confirmed_block::RewardType::DeactivatedStake as i32 => { + crate::YellowstoneRewardType::DeactivatedStake + }, + value => crate::YellowstoneRewardType::Unknown(value), + }; + return std::result::Result::Ok(crate::YellowstoneReward { + pubkey, + lamports: wire.lamports, + post_balance: wire.post_balance, + reward_type, + commission: wire.commission, + commission_bps: wire.commission_bps, + }); +} + +#[cfg(test)] +fn decode_transaction_error( + field: &'static str, + wire: yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionError, +) -> ksp_core_lib::Result { + if wire.err.len() > MAX_GRPC_TRANSACTION_ERROR_BYTES { + return invalid_subscribe_response(field, "Yellowstone transaction error bytes exceed the KSP bound"); + } + return std::result::Result::Ok(crate::YellowstoneTransactionError { bytes: wire.err }); +} + +#[cfg(test)] +fn decode_transaction_signature(field: &'static str, bytes: std::vec::Vec) -> ksp_core_lib::Result { + let bytes: [u8; YELLOWSTONE_TRANSACTION_SIGNATURE_WIRE_LENGTH_BYTES] = match bytes.try_into() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return invalid_subscribe_response(field, "Yellowstone transaction contains an invalid signature length"), + }; + return std::result::Result::Ok(crate::YellowstoneTransactionSignature::new(bytes)); +} + +#[cfg(test)] +fn decode_hash_bytes(field: &'static str, bytes: std::vec::Vec) -> ksp_core_lib::Result { + let bytes: [u8; YELLOWSTONE_HASH_LENGTH_BYTES] = match bytes.try_into() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return invalid_subscribe_response(field, "Yellowstone transaction contains an invalid hash length"), + }; + return std::result::Result::Ok(crate::YellowstoneHashBytes { bytes }); +} + #[cfg(test)] fn decode_update_envelope( filters: std::vec::Vec, @@ -1322,7 +2842,7 @@ fn decode_account_info(wire: yellowstone_grpc_proto::geyser::SubscribeUpdateAcco }; let transaction_signature = match wire.txn_signature { std::option::Option::Some(value) => { - let bytes: [u8; YELLOWSTONE_TRANSACTION_SIGNATURE_LENGTH_BYTES] = match value.try_into() { + let bytes: [u8; YELLOWSTONE_TRANSACTION_SIGNATURE_WIRE_LENGTH_BYTES] = match value.try_into() { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => { return invalid_subscribe_response("account.txn_signature", "Yellowstone account update contains an invalid transaction signature length"); @@ -1353,6 +2873,59 @@ fn decode_pubkey_bytes(field: &'static str, bytes: std::vec::Vec) -> ksp_cor return std::result::Result::Ok(ksp_core_lib::Pubkey::new_from_array(array)); } +fn push_transaction_selector(target: &mut std::vec::Vec, value: ksp_core_lib::Pubkey, field: &'static str) -> ksp_core_lib::Result<()> { + if target.len() >= MAX_GRPC_SUBSCRIBE_ACCOUNT_SELECTOR_COUNT { + return std::result::Result::Err( + ksp_core_lib::Error::new(crate::ERROR_CODE_INVALID_RPC_PARAMETERS, "Yellowstone transaction account selector count exceeds KSP bounds") + .with_context("field", field), + ); + } + target.push(value); + return std::result::Result::Ok(()); +} + +fn base58_decoded_length(value: &str) -> std::option::Option { + let leading_zero_bytes = value.bytes().take_while(|value| return *value == b'1').count(); + let mut decoded = std::vec::Vec::::new(); + for value in value.bytes() { + let digit = match base58_digit(value) { + std::option::Option::Some(value) => u32::from(value), + std::option::Option::None => return std::option::Option::None, + }; + let mut carry = digit; + for byte in &mut decoded { + let expanded = u32::from(*byte) * 58 + carry; + let low = match u8::try_from(expanded & 0xff) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::option::Option::None, + }; + *byte = low; + carry = expanded >> 8; + } + while carry > 0 { + let low = match u8::try_from(carry & 0xff) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::option::Option::None, + }; + decoded.push(low); + carry >>= 8; + } + } + return leading_zero_bytes.checked_add(decoded.len()); +} + +const fn base58_digit(value: u8) -> std::option::Option { + return match value { + b'1'..=b'9' => std::option::Option::Some(value - b'1'), + b'A'..=b'H' => std::option::Option::Some(value - b'A' + 9), + b'J'..=b'N' => std::option::Option::Some(value - b'J' + 17), + b'P'..=b'Z' => std::option::Option::Some(value - b'P' + 22), + b'a'..=b'k' => std::option::Option::Some(value - b'a' + 33), + b'm'..=b'z' => std::option::Option::Some(value - b'm' + 44), + _ => std::option::Option::None, + }; +} + fn validate_memcmp_text(value: &str) -> ksp_core_lib::Result<()> { if value.len() > MAX_GRPC_SUBSCRIBE_MEMCMP_TEXT_LENGTH_BYTES || !value.is_ascii() || value.chars().any(char::is_whitespace) { return invalid_subscribe_parameter("grpc_subscribe.accounts.filters.memcmp.text", "Yellowstone memcmp text violates KSP bounds"); diff --git a/crates/ksp-onchain-transport-lib/src/lib.rs b/crates/ksp-onchain-transport-lib/src/lib.rs index 77f17d4..f7faa08 100644 --- a/crates/ksp-onchain-transport-lib/src/lib.rs +++ b/crates/ksp-onchain-transport-lib/src/lib.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/lib.rs -// version: 39 +// version: 40 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -36,8 +36,10 @@ //! `0.2.9-pre.003` adds bounded TLS/WebPKI connection establishment, generic redacted ASCII request metadata and the seven standard Yellowstone unary RPCs //! through KSP-owned DTOs. //! `0.2.9-pre.004` materializes the provider-neutral standard `SubscribeRequest` foundation: all seven named filter maps, global filter-name bounds/uniqueness, -//! commitment, ordered account-data slices, ping and `from_slot`. Family-specific account/slot filters land in `pre.005`; transaction/block filters remain staged for `pre.007–008`. -//! `0.2.9-pre.006` normalizes the five unambiguously HTTP-owned private implementation modules with an `http_` prefix while preserving shared `rpc_*`, JSON-RPC, error and constants modules. +//! commitment, ordered account-data slices, ping and `from_slot`. Family-specific account/slot filters land in `pre.005`; transaction/block filters remain +//! staged for `pre.007–008`. +//! `0.2.9-pre.006` normalizes the five unambiguously HTTP-owned private implementation modules with an `http_` prefix while preserving shared `rpc_*`, +//! JSON-RPC, error and constants modules. mod constants; mod error; @@ -71,14 +73,6 @@ mod ws_settings; mod ws_subscription; mod ws_transactions; -/// Passive runtime availability reported for one logical HTTP endpoint. -pub use self::http_client::HttpEndpointAvailability; -/// Shareable logical HTTP endpoint client owned by KSP Transport. -pub use self::http_client::HttpEndpointClient; -/// Safe routing snapshot for one configured endpoint role. -pub use self::http_client::HttpEndpointRoleSnapshot; -/// Safe metadata snapshot for one logical HTTP endpoint. -pub use self::http_client::HttpEndpointSnapshot; /// Error code used when no logical endpoint can satisfy a request. pub use self::error::ERROR_CODE_ENDPOINT_SELECTION_FAILED; /// Error code used when a Yellowstone gRPC channel cannot be prepared safely. @@ -135,24 +129,46 @@ pub use self::grpc_settings::YellowstoneGrpcReconnectSettings; pub use self::grpc_settings::YellowstoneGrpcSessionSettings; /// Complete runtime settings consumed by the KSP Yellowstone gRPC transport engine. pub use self::grpc_settings::YellowstoneGrpcTransportSettings; -/// Typed account payload carried by one standard Yellowstone account update. -pub use self::grpc_subscribe::YellowstoneAccountInfo; /// One validated standard Yellowstone account predicate. pub use self::grpc_subscribe::YellowstoneAccountFilterPredicate; -/// Standard Yellowstone account-update projection owned by KSP. -pub use self::grpc_subscribe::YellowstoneAccountUpdate; +/// Typed account payload carried by one standard Yellowstone account update. +pub use self::grpc_subscribe::YellowstoneAccountInfo; +/// Lamport comparison used by standard Yellowstone account filters. +pub use self::grpc_subscribe::YellowstoneAccountLamportsFilter; /// Validated standard Yellowstone account memcmp predicate. pub use self::grpc_subscribe::YellowstoneAccountMemcmp; /// Encoding selected by one Yellowstone account memcmp predicate. pub use self::grpc_subscribe::YellowstoneAccountMemcmpEncoding; -/// Lamport comparison used by standard Yellowstone account filters. -pub use self::grpc_subscribe::YellowstoneAccountLamportsFilter; +/// Standard Yellowstone account-update projection owned by KSP. +pub use self::grpc_subscribe::YellowstoneAccountUpdate; /// One standard Yellowstone account-data slice. pub use self::grpc_subscribe::YellowstoneAccountsDataSlice; -/// Hash algorithm carried by a standard Yellowstone Cuckoo filter. -pub use self::grpc_subscribe::YellowstoneCuckooHashAlgorithm; +/// One compiled instruction from the Yellowstone Solana-storage transaction wire. +pub use self::grpc_subscribe::YellowstoneCompiledInstruction; /// Wire-preserving KSP representation of a standard Yellowstone Cuckoo filter. pub use self::grpc_subscribe::YellowstoneCuckooFilter; +/// Hash algorithm carried by a standard Yellowstone Cuckoo filter. +pub use self::grpc_subscribe::YellowstoneCuckooHashAlgorithm; +/// Fixed-width 32-byte hash from the Yellowstone Solana-storage transaction wire. +pub use self::grpc_subscribe::YellowstoneHashBytes; +/// One inner instruction from Yellowstone transaction status metadata. +pub use self::grpc_subscribe::YellowstoneInnerInstruction; +/// One indexed inner-instruction group from Yellowstone transaction status metadata. +pub use self::grpc_subscribe::YellowstoneInnerInstructions; +/// One address-table lookup from a Yellowstone transaction message. +pub use self::grpc_subscribe::YellowstoneMessageAddressTableLookup; +/// Return-data payload from Yellowstone transaction status metadata. +pub use self::grpc_subscribe::YellowstoneReturnData; +/// One reward entry from Yellowstone transaction status metadata. +pub use self::grpc_subscribe::YellowstoneReward; +/// Reward classification from the Yellowstone Solana-storage wire. +pub use self::grpc_subscribe::YellowstoneRewardType; +/// Current standard Yellowstone slot status. +pub use self::grpc_subscribe::YellowstoneSlotStatus; +/// Standard Yellowstone slot-update projection owned by KSP. +pub use self::grpc_subscribe::YellowstoneSlotUpdate; +/// Solana transaction body carried by Yellowstone storage protobuf messages. +pub use self::grpc_subscribe::YellowstoneStoredTransaction; /// Complete account-family filter group for standard Yellowstone Subscribe. pub use self::grpc_subscribe::YellowstoneSubscribeAccountFilter; /// Block-family filter-group shell for standard Yellowstone Subscribe. @@ -169,16 +185,36 @@ pub use self::grpc_subscribe::YellowstoneSubscribePing; pub use self::grpc_subscribe::YellowstoneSubscribeRequest; /// Complete slot-family filter group for standard Yellowstone Subscribe. pub use self::grpc_subscribe::YellowstoneSubscribeSlotFilter; -/// Current standard Yellowstone slot status. -pub use self::grpc_subscribe::YellowstoneSlotStatus; -/// Standard Yellowstone slot-update projection owned by KSP. -pub use self::grpc_subscribe::YellowstoneSlotUpdate; -/// Fixed-width transaction signature attached to Yellowstone account updates when available. +/// Complete transaction-family filter shared by transactions and transaction-status maps. +pub use self::grpc_subscribe::YellowstoneSubscribeTransactionFilter; +/// Optional token-account owner expansion for current Yellowstone transaction filters. +pub use self::grpc_subscribe::YellowstoneTokenAccountExpansion; +/// One pre/post token balance from Yellowstone transaction status metadata. +pub use self::grpc_subscribe::YellowstoneTokenBalance; +/// Optional Transaction V1 inline budget configuration from Yellowstone Solana-storage. +pub use self::grpc_subscribe::YellowstoneTransactionConfig; +/// Opaque runtime transaction error bytes from Yellowstone Solana-storage. +pub use self::grpc_subscribe::YellowstoneTransactionError; +/// Complete transaction info carried by Yellowstone transaction and block updates. +pub use self::grpc_subscribe::YellowstoneTransactionInfo; +/// Complete current Yellowstone transaction message. +pub use self::grpc_subscribe::YellowstoneTransactionMessage; +/// Solana transaction message header from Yellowstone Solana-storage. +pub use self::grpc_subscribe::YellowstoneTransactionMessageHeader; +/// Fixed-width transaction signature attached to Yellowstone updates. pub use self::grpc_subscribe::YellowstoneTransactionSignature; +/// Validated base58 transaction-signature selector for Yellowstone transaction filters. +pub use self::grpc_subscribe::YellowstoneTransactionSignatureSelector; +/// Complete Yellowstone transaction status metadata. +pub use self::grpc_subscribe::YellowstoneTransactionStatusMeta; +/// Lightweight Yellowstone transaction-status update. +pub use self::grpc_subscribe::YellowstoneTransactionStatusUpdate; +/// Full Yellowstone transaction update. +pub use self::grpc_subscribe::YellowstoneTransactionUpdate; +/// UI token amount from Yellowstone transaction status metadata. +pub use self::grpc_subscribe::YellowstoneUiTokenAmount; /// Timestamp attached to standard Yellowstone update envelopes. pub use self::grpc_subscribe::YellowstoneUpdateTimestamp; -/// Transaction-family filter-group shell shared by transactions and transaction-status maps. -pub use self::grpc_subscribe::YellowstoneSubscribeTransactionFilter; /// Standard Solana Yellowstone unary facade over one KSP-owned physical gRPC channel. pub use self::grpc_unary::SolanaYellowstoneGrpcUnaryClient; /// Block height returned by the standard Yellowstone unary surface. @@ -195,20 +231,14 @@ pub use self::grpc_unary::YellowstoneReplayInfo; pub use self::grpc_unary::YellowstoneSlot; /// Bounded endpoint version returned by the standard Yellowstone unary surface. pub use self::grpc_unary::YellowstoneVersionInfo; -/// JSON-RPC 2.0 error payload returned by a remote Solana endpoint. -pub use self::json_rpc::JsonRpcErrorObject; -/// Validated JSON-RPC 2.0 error response. -pub use self::json_rpc::JsonRpcErrorResponse; -/// JSON-RPC 2.0 HTTP request envelope emitted by KSP. -pub use self::json_rpc::JsonRpcRequest; -/// Validated JSON-RPC 2.0 HTTP response. -pub use self::json_rpc::JsonRpcResponse; -/// Validated JSON-RPC 2.0 success response. -pub use self::json_rpc::JsonRpcSuccessResponse; -/// Parses and validates a JSON-RPC HTTP response from UTF-8 JSON text. -pub use self::json_rpc::parse_json_rpc_response_text; -/// Validates a decoded JSON value as one JSON-RPC HTTP response. -pub use self::json_rpc::parse_json_rpc_response_value; +/// Passive runtime availability reported for one logical HTTP endpoint. +pub use self::http_client::HttpEndpointAvailability; +/// Shareable logical HTTP endpoint client owned by KSP Transport. +pub use self::http_client::HttpEndpointClient; +/// Safe routing snapshot for one configured endpoint role. +pub use self::http_client::HttpEndpointRoleSnapshot; +/// Safe metadata snapshot for one logical HTTP endpoint. +pub use self::http_client::HttpEndpointSnapshot; /// Result of one logical endpoint selection. pub use self::http_pool::HttpEndpointSelection; /// Runtime admission permit for one HTTP request. @@ -225,6 +255,40 @@ pub use self::http_resilience::HttpRetryCause; pub use self::http_resilience::HttpRetryDecision; /// Evaluates the centralized bounded HTTP retry policy for one audited RPC method. pub use self::http_resilience::evaluate_transport_retry; +/// Open cluster or network descriptor used by HTTP endpoint settings. +pub use self::http_settings::HttpClusterName; +/// Runtime settings for one role declared by an HTTP endpoint. +pub use self::http_settings::HttpEndpointRoleSettings; +/// Runtime settings for one named Solana HTTP endpoint. +pub use self::http_settings::HttpEndpointSettings; +/// Runtime HTTP endpoint URL with redacted diagnostics. +pub use self::http_settings::HttpEndpointUrl; +/// Open provider descriptor used by HTTP endpoint settings. +pub use self::http_settings::HttpProviderName; +/// Open request-kind descriptor used by logical endpoint capabilities. +pub use self::http_settings::HttpRequestKind; +/// Bounded retry settings owned by the HTTP transport runtime. +pub use self::http_settings::HttpRetrySettings; +/// Local limits attached to one logical HTTP endpoint role. +pub use self::http_settings::HttpRoleLimits; +/// Open logical endpoint role descriptor. +pub use self::http_settings::HttpRoleName; +/// Complete runtime settings consumed by the Solana HTTP transport foundation. +pub use self::http_settings::HttpTransportSettings; +/// JSON-RPC 2.0 error payload returned by a remote Solana endpoint. +pub use self::json_rpc::JsonRpcErrorObject; +/// Validated JSON-RPC 2.0 error response. +pub use self::json_rpc::JsonRpcErrorResponse; +/// JSON-RPC 2.0 HTTP request envelope emitted by KSP. +pub use self::json_rpc::JsonRpcRequest; +/// Validated JSON-RPC 2.0 HTTP response. +pub use self::json_rpc::JsonRpcResponse; +/// Validated JSON-RPC 2.0 success response. +pub use self::json_rpc::JsonRpcSuccessResponse; +/// Parses and validates a JSON-RPC HTTP response from UTF-8 JSON text. +pub use self::json_rpc::parse_json_rpc_response_text; +/// Validates a decoded JSON value as one JSON-RPC HTTP response. +pub use self::json_rpc::parse_json_rpc_response_value; /// Typed transport-level Solana account without Program/SPL decoding. pub use self::rpc_accounts::SolanaAccount; /// Address and lamport balance returned by `getLargestAccounts`. @@ -397,26 +461,6 @@ pub use self::rpc_transactions::SolanaTransactionEncoding; pub use self::rpc_transactions::SolanaTransactionVersion; /// Three-state wire field used when Solana distinguishes omission from an explicit JSON `null`. pub use self::rpc_transactions::SolanaWireField; -/// Open cluster or network descriptor used by HTTP endpoint settings. -pub use self::http_settings::HttpClusterName; -/// Runtime settings for one role declared by an HTTP endpoint. -pub use self::http_settings::HttpEndpointRoleSettings; -/// Runtime settings for one named Solana HTTP endpoint. -pub use self::http_settings::HttpEndpointSettings; -/// Runtime HTTP endpoint URL with redacted diagnostics. -pub use self::http_settings::HttpEndpointUrl; -/// Open provider descriptor used by HTTP endpoint settings. -pub use self::http_settings::HttpProviderName; -/// Open request-kind descriptor used by logical endpoint capabilities. -pub use self::http_settings::HttpRequestKind; -/// Bounded retry settings owned by the HTTP transport runtime. -pub use self::http_settings::HttpRetrySettings; -/// Local limits attached to one logical HTTP endpoint role. -pub use self::http_settings::HttpRoleLimits; -/// Open logical endpoint role descriptor. -pub use self::http_settings::HttpRoleName; -/// Complete runtime settings consumed by the Solana HTTP transport foundation. -pub use self::http_settings::HttpTransportSettings; /// Configuration accepted by the standard Solana `accountSubscribe` WebSocket method. pub use self::ws_accounts::SolanaAccountSubscribeConfig; /// One `programNotification` payload preserving contextual and non-contextual upstream forms. @@ -510,12 +554,12 @@ pub(crate) use self::http_resilience::HttpConcurrencyPermit; pub(crate) use self::http_resilience::HttpRoleRuntime; /// Crate-internal `RoleAdmissionAttempt` variants used by the owning crate. pub(crate) use self::http_resilience::RoleAdmissionAttempt; +/// Validates endpoint settings. +pub(crate) use self::http_settings::validate_endpoint_settings; /// Decodes one private serde wire type into the shared Transport error domain for typed RPC adapters. pub(crate) use self::rpc_common::decode_wire_json; /// Parses a base58 public key without echoing its wire value into diagnostics for typed RPC adapters. pub(crate) use self::rpc_common::parse_wire_pubkey; -/// Validates endpoint settings. -pub(crate) use self::http_settings::validate_endpoint_settings; /// Crate-internal command surface shared by the physical session and typed subscription handle. pub(crate) use self::ws_session::WsSessionCommand; /// Crate-internal notification dispatch result. diff --git a/crates/ksp-onchain-transport-lib/tests/public_api.rs b/crates/ksp-onchain-transport-lib/tests/public_api.rs index a551a06..6bf0ca1 100644 --- a/crates/ksp-onchain-transport-lib/tests/public_api.rs +++ b/crates/ksp-onchain-transport-lib/tests/public_api.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/tests/public_api.rs -// version: 43 +// version: 44 //! Integration tests for the public `ksp-onchain-transport-lib` consumer contract. @@ -896,3 +896,32 @@ fn public_v0_2_9_pre_005_yellowstone_accounts_slots_contract_is_available_from_c let _lamports = ksp_onchain_transport_lib::YellowstoneAccountLamportsFilter::Gt(10); let _encoding = ksp_onchain_transport_lib::YellowstoneAccountMemcmpEncoding::Bytes; } + +#[test] +fn public_v0_2_9_pre_007_yellowstone_transactions_contract_is_available_from_crate_root() { + let mut filter = ksp_onchain_transport_lib::YellowstoneSubscribeTransactionFilter::new(); + filter.set_vote(std::option::Option::Some(false)); + filter.set_failed(std::option::Option::Some(true)); + let selector = ksp_onchain_transport_lib::YellowstoneTransactionSignatureSelector::new("1".repeat(64)).expect("signature selector must validate"); + filter.set_signature(std::option::Option::Some(selector)); + assert!(filter.push_account_include(ksp_core_lib::Pubkey::new_from_array([1_u8; 32])).is_ok()); + assert!(filter.push_account_exclude(ksp_core_lib::Pubkey::new_from_array([2_u8; 32])).is_ok()); + assert!(filter.push_account_required(ksp_core_lib::Pubkey::new_from_array([3_u8; 32])).is_ok()); + filter.set_token_accounts(std::option::Option::Some(ksp_onchain_transport_lib::YellowstoneTokenAccountExpansion::All)); + assert_eq!(filter.token_accounts(), std::option::Option::Some(ksp_onchain_transport_lib::YellowstoneTokenAccountExpansion::All)); + let _hash = std::any::type_name::(); + let _instruction = std::any::type_name::(); + let _lookup = std::any::type_name::(); + let _config = std::any::type_name::(); + let _message = std::any::type_name::(); + let _stored = std::any::type_name::(); + let _meta = std::any::type_name::(); + let _info = std::any::type_name::(); + let _update = std::any::type_name::(); + let _status = std::any::type_name::(); + let _error = std::any::type_name::(); + let _inner = std::any::type_name::(); + let _token = std::any::type_name::(); + let _return_data = std::any::type_name::(); + let _reward = std::any::type_name::(); +} diff --git a/crates/ksp-onchain-transport-lib/tests/release_completeness.rs b/crates/ksp-onchain-transport-lib/tests/release_completeness.rs index 7ed02c0..2088792 100644 --- a/crates/ksp-onchain-transport-lib/tests/release_completeness.rs +++ b/crates/ksp-onchain-transport-lib/tests/release_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/tests/release_completeness.rs -// version: 37 +// version: 38 //! Release-level completeness canaries for staged HTTP and WebSocket Transport coverage. @@ -1071,7 +1071,19 @@ fn release_v0_2_9_pre_003_adds_tls_metadata_and_exactly_seven_standard_unary_met fn release_v0_2_9_pre_004_materializes_only_standard_subscribe_common_contract() { let source = include_str!("../src/grpc_subscribe.rs"); let crate_root = include_str!("../src/lib.rs"); - for wire_field in ["accounts:", "slots:", "transactions:", "transactions_status:", "blocks:", "blocks_meta:", "entry:", "commitment:", "accounts_data_slice:", "ping:", "from_slot:"] { + for wire_field in [ + "accounts:", + "slots:", + "transactions:", + "transactions_status:", + "blocks:", + "blocks_meta:", + "entry:", + "commitment:", + "accounts_data_slice:", + "ping:", + "from_slot:", + ] { assert!(source.contains(wire_field), "missing standard Yellowstone SubscribeRequest field: {wire_field}"); } assert!(source.contains("MAX_GRPC_SUBSCRIBE_FILTER_GROUP_COUNT")); @@ -1134,7 +1146,6 @@ fn release_v0_2_9_pre_005_completes_standard_accounts_and_slots_without_advancin assert!(source.contains("MAX_GRPC_SUBSCRIBE_SLOT_DEAD_ERROR_LENGTH_BYTES")); assert!(source.contains("#[cfg(test)]\nfn decode_account_update")); assert!(source.contains("#[cfg(test)]\nfn decode_slot_update")); - assert!(!source.contains("YellowstoneTransactionUpdate")); assert!(!source.contains("YellowstoneBlockUpdate")); assert!(!source.contains("SubscribeDeshred")); assert!(!source.contains("PublicNode")); @@ -1168,3 +1179,45 @@ fn release_v0_2_9_pre_006_namespaces_unambiguously_http_owned_private_modules() assert!(crate_root.contains(shared), "shared/protocol module was incorrectly HTTP-prefixed: {shared}"); } } + +#[test] +fn release_v0_2_9_pre_007_completes_standard_transactions_without_advancing_blocks_or_bidi() { + let source = include_str!("../src/grpc_subscribe.rs"); + let crate_root = include_str!("../src/lib.rs"); + for required in [ + "YellowstoneSubscribeTransactionFilter", + "YellowstoneTransactionSignatureSelector", + "YellowstoneTokenAccountExpansion", + "account_include", + "account_exclude", + "account_required", + "cuckoo_account_include", + "token_accounts", + "YellowstoneTransactionUpdate", + "YellowstoneTransactionStatusUpdate", + "YellowstoneTransactionStatusMeta", + "YellowstoneTransactionMessage", + "YellowstoneMessageAddressTableLookup", + "YellowstoneTransactionConfig", + "compute_units_consumed", + "cost_units", + "inner_instructions_none", + "log_messages_none", + "return_data_none", + "decode_transaction_update", + "decode_transaction_status_update", + ] { + assert!(source.contains(required), "missing pre.007 Transactions contract token: {required}"); + } + assert!(source.contains("base58_decoded_length")); + assert!(source.contains("YELLOWSTONE_TRANSACTION_SIGNATURE_WIRE_LENGTH_BYTES")); + assert!(!source.contains("YellowstoneBlockUpdate")); + assert!(!source.contains("fn decode_block_update")); + assert!(!source.contains("SubscribeDeshred")); + assert!(!source.contains("PublicNode")); + assert!(!crate_root.contains("pub use tonic")); + assert!(!crate_root.contains("pub use yellowstone_grpc_proto")); + let _filter = std::any::type_name::(); + let _update = std::any::type_name::(); + let _status = std::any::type_name::(); +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/client.rs b/crates/ksp-onchain-transport-lib/unit_tests/client.rs deleted file mode 100644 index 700f04c..0000000 --- a/crates/ksp-onchain-transport-lib/unit_tests/client.rs +++ /dev/null @@ -1,62 +0,0 @@ -// file: crates/ksp-onchain-transport-lib/unit_tests/client.rs -// version: 3 - -fn endpoint(enabled: bool, url_text: &str) -> crate::HttpEndpointSettings { - let url = crate::HttpEndpointUrl::parse(url_text).expect("test endpoint URL must parse"); - let role = crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new("default"), - true, - std::vec![crate::HttpRequestKind::wildcard()], - 10, - crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), - ); - return crate::HttpEndpointSettings::new( - "endpoint", - enabled, - crate::HttpProviderName::new("provider"), - crate::HttpClusterName::new("devnet"), - url, - std::time::Duration::from_secs(1), - std::time::Duration::from_secs(2), - std::option::Option::Some(4), - std::vec![role], - ); -} - -#[test] -fn endpoint_client_snapshot_never_contains_url_or_secret_material() { - let client = crate::HttpEndpointClient::new(endpoint(true, "https://provider.invalid/rpc?api-key=SECRET-CANARY")).expect("client must build"); - let snapshot = client.snapshot(); - let rendered = format!("{snapshot:?} {client:?}"); - assert_eq!(snapshot.availability(), crate::HttpEndpointAvailability::Available); - assert!(!rendered.contains("SECRET-CANARY")); - assert!(!rendered.contains("provider.invalid")); - assert!(!rendered.contains("https://")); -} - -#[test] -fn disabled_endpoint_client_is_visible_but_not_selectable() { - let client = crate::HttpEndpointClient::new(endpoint(false, "https://api.devnet.solana.com")).expect("disabled client must still build"); - assert_eq!(client.snapshot().availability(), crate::HttpEndpointAvailability::Disabled); - assert!(!client.supports(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance"))); -} - -#[test] -fn endpoint_client_matches_exact_and_wildcard_capabilities() { - let client = crate::HttpEndpointClient::new(endpoint(true, "https://api.devnet.solana.com")).expect("client must build"); - assert!(client.supports(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance"))); - assert!(!client.supports(&crate::HttpRoleName::new("write"), &crate::HttpRequestKind::new("get_balance"))); -} - -#[test] -fn endpoint_role_snapshot_exposes_safe_resilience_state() { - let client = crate::HttpEndpointClient::new(endpoint(true, "https://api.devnet.solana.com")).expect("client must build"); - let snapshot = client.snapshot(); - let role = &snapshot.roles()[0]; - assert_eq!(role.availability(), crate::HttpEndpointAvailability::Available); - assert_eq!(role.in_flight_requests(), std::option::Option::None); - assert_eq!(role.cooldown_remaining(), std::option::Option::None); - assert_eq!(role.success_count(), 0); - assert_eq!(role.failure_count(), 0); - assert_eq!(role.rate_limit_count(), 0); -} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/executor.rs b/crates/ksp-onchain-transport-lib/unit_tests/executor.rs deleted file mode 100644 index c6419c9..0000000 --- a/crates/ksp-onchain-transport-lib/unit_tests/executor.rs +++ /dev/null @@ -1,141 +0,0 @@ -// file: crates/ksp-onchain-transport-lib/unit_tests/executor.rs -// version: 2 - -fn pool_for_url(url: &str, request_timeout: std::time::Duration, max_retries: u32) -> crate::HttpTransportPool { - let role = crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new("default"), - true, - std::vec![crate::HttpRequestKind::wildcard()], - 10, - crate::HttpRoleLimits::new( - std::option::Option::None, - std::option::Option::None, - std::option::Option::None, - std::option::Option::Some(std::time::Duration::from_millis(1)), - ), - ); - let endpoint = crate::HttpEndpointSettings::new( - "fixture", - true, - crate::HttpProviderName::new("fixture"), - crate::HttpClusterName::new("local"), - crate::HttpEndpointUrl::parse(url).expect("fixture URL must parse"), - std::time::Duration::from_millis(100), - request_timeout, - std::option::Option::Some(1), - std::vec![role], - ); - let settings = crate::HttpTransportSettings::new( - std::vec![endpoint], - crate::HttpRetrySettings::new(max_retries, std::time::Duration::from_millis(1), std::time::Duration::from_millis(2)), - ); - return crate::HttpTransportPool::new(settings).expect("fixture pool must build"); -} - -fn health_method() -> &'static crate::HttpRpcMethodDescriptor { - return crate::find_http_rpc_method("getHealth").expect("getHealth descriptor must exist"); -} - -fn serve_rate_limit_then_success() -> (std::string::String, std::thread::JoinHandle) { - let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("fixture listener must bind"); - let address = listener.local_addr().expect("fixture listener address must resolve"); - let handle = std::thread::spawn(move || { - let mut count = 0_usize; - while count < 2 { - let (mut stream, _) = listener.accept().expect("fixture server must accept request"); - let _ = read_request(&mut stream); - let response = if count == 0 { - "HTTP/1.1 429 Too Many Requests\r\nRetry-After: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n".to_owned() - } else { - let body = include_str!("../fixtures/http/get_health.success.json"); - format!("HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", body.len(), body) - }; - std::io::Write::write_all(&mut stream, response.as_bytes()).expect("fixture response must write"); - count = count.saturating_add(1); - } - return count; - }); - return (format!("http://{address}"), handle); -} - -fn serve_timeout() -> (std::string::String, std::thread::JoinHandle<()>) { - let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("fixture listener must bind"); - let address = listener.local_addr().expect("fixture listener address must resolve"); - let handle = std::thread::spawn(move || { - let (mut stream, _) = listener.accept().expect("fixture server must accept request"); - let _ = read_request(&mut stream); - std::thread::sleep(std::time::Duration::from_millis(100)); - return; - }); - return (format!("http://{address}"), handle); -} - -fn read_request(stream: &mut std::net::TcpStream) -> std::string::String { - let mut bytes = std::vec::Vec::new(); - let mut buffer = [0_u8; 1024]; - loop { - let count = std::io::Read::read(stream, &mut buffer).expect("fixture request must read"); - if count == 0 { - break; - } - bytes.extend_from_slice(&buffer[..count]); - if request_complete(bytes.as_slice()) { - break; - } - } - return std::string::String::from_utf8(bytes).expect("fixture request must be UTF-8"); -} - -fn request_complete(bytes: &[u8]) -> bool { - let text = match std::str::from_utf8(bytes) { - std::result::Result::Ok(text) => text, - std::result::Result::Err(_) => return false, - }; - let header_end = match text.find("\r\n\r\n") { - std::option::Option::Some(value) => value, - std::option::Option::None => return false, - }; - let mut content_length = 0_usize; - for line in text[..header_end].lines() { - let (name, value) = match line.split_once(':') { - std::option::Option::Some(parts) => parts, - std::option::Option::None => continue, - }; - if name.eq_ignore_ascii_case("content-length") { - content_length = value.trim().parse::().expect("content length must parse"); - } - } - return bytes.len() >= header_end.saturating_add(4).saturating_add(content_length); -} - -#[tokio::test(flavor = "current_thread")] -async fn executor_applies_retry_after_and_retries_http_429_for_retry_safe_method() { - let (url, handle) = serve_rate_limit_then_success(); - let pool = pool_for_url(url.as_str(), std::time::Duration::from_millis(500), 1); - let result = pool - .execute_standard_rpc(&crate::HttpRoleName::new("default"), health_method(), std::vec::Vec::new()) - .await - .expect("retry-safe request must recover from one 429"); - assert_eq!(result, serde_json::json!("ok")); - assert_eq!(handle.join().expect("fixture server must join"), 2); - let snapshot = pool.snapshot(); - assert_eq!(snapshot.endpoints()[0].roles()[0].rate_limit_count(), 1); - assert_eq!(snapshot.endpoints()[0].roles()[0].success_count(), 1); -} - -#[tokio::test(flavor = "current_thread")] -async fn executor_maps_reqwest_timeout_to_ksp_timeout_error_without_endpoint_secret_leak() { - const SECRET_CANARY: &str = "SECRET-REQWEST-URL-CANARY"; - let (url, handle) = serve_timeout(); - let sensitive_url = format!("{url}/rpc?api-key={SECRET_CANARY}"); - let pool = pool_for_url(sensitive_url.as_str(), std::time::Duration::from_millis(20), 0); - let error = pool - .execute_standard_rpc(&crate::HttpRoleName::new("default"), health_method(), std::vec::Vec::new()) - .await - .expect_err("timed out request must fail"); - assert_eq!(error.code(), crate::ERROR_CODE_TIMEOUT); - assert!(!format!("{error:?}").contains(SECRET_CANARY)); - let source = std::error::Error::source(&error).expect("transport timeout should preserve a sanitized reqwest source"); - assert!(!format!("{source:?}").contains(SECRET_CANARY)); - handle.join().expect("fixture server must join"); -} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/grpc_subscribe.rs b/crates/ksp-onchain-transport-lib/unit_tests/grpc_subscribe.rs index fa343b6..6c6157b 100644 --- a/crates/ksp-onchain-transport-lib/unit_tests/grpc_subscribe.rs +++ b/crates/ksp-onchain-transport-lib/unit_tests/grpc_subscribe.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/unit_tests/grpc_subscribe.rs -// version: 2 +// version: 3 fn filter_name(value: &str) -> crate::YellowstoneSubscribeFilterName { return crate::YellowstoneSubscribeFilterName::new(value).expect("fixture filter name must validate"); @@ -55,7 +55,7 @@ fn yellowstone_subscribe_empty_and_named_empty_maps_encode_exactly() { assert_eq!(wire.slots.get("slots"), std::option::Option::Some(&yellowstone_grpc_proto::geyser::SubscribeRequestFilterSlots::default())); assert_eq!( wire.transactions.get("transactions"), - std::option::Option::Some(&yellowstone_grpc_proto::geyser::SubscribeRequestFilterTransactions::default()) + std::option::Option::Some(&yellowstone_grpc_proto::geyser::SubscribeRequestFilterTransactions::default()), ); assert_eq!( wire.transactions_status.get("transaction-status"), @@ -365,3 +365,191 @@ fn yellowstone_slot_update_preserves_all_current_statuses_and_bounds_dead_error( }; assert!(super::decode_slot_update(oversized).is_err()); } + +#[test] +fn yellowstone_transaction_filters_encode_complete_current_wire_and_redact_selectors() { + let mut filter = crate::YellowstoneSubscribeTransactionFilter::new(); + filter.set_vote(std::option::Option::Some(false)); + filter.set_failed(std::option::Option::Some(true)); + let signature = crate::YellowstoneTransactionSignatureSelector::new("1".repeat(64)).expect("signature selector must validate"); + filter.set_signature(std::option::Option::Some(signature)); + assert!(filter.push_account_include(ksp_core_lib::Pubkey::new_from_array([1_u8; 32])).is_ok()); + assert!(filter.push_account_exclude(ksp_core_lib::Pubkey::new_from_array([2_u8; 32])).is_ok()); + assert!(filter.push_account_required(ksp_core_lib::Pubkey::new_from_array([3_u8; 32])).is_ok()); + let cuckoo = + crate::YellowstoneCuckooFilter::new(vec![0_u8; 16], 4, 4, 8, 7, crate::YellowstoneCuckooHashAlgorithm::SipHash).expect("cuckoo filter must validate"); + filter.set_cuckoo_account_include(std::option::Option::Some(cuckoo)); + filter.set_token_accounts(std::option::Option::Some(crate::YellowstoneTokenAccountExpansion::BalanceChanged)); + let wire = filter.to_wire(); + assert_eq!(wire.vote, std::option::Option::Some(false)); + assert_eq!(wire.failed, std::option::Option::Some(true)); + assert!(wire.signature.is_some()); + assert_eq!(wire.account_include.len(), 1); + assert_eq!(wire.account_exclude.len(), 1); + assert_eq!(wire.account_required.len(), 1); + assert!(wire.cuckoo_account_include.is_some()); + assert_eq!(wire.token_accounts, std::option::Option::Some(yellowstone_grpc_proto::geyser::TokenAccountExpansionControlFlag::BalanceChanged as i32)); + let debug = format!("{filter:?}"); + assert!(debug.contains("account_include_count")); + assert!(!debug.contains(&"1".repeat(32))); + assert!(!debug.contains(&ksp_core_lib::Pubkey::new_from_array([1_u8; 32]).to_string())); + assert!(crate::YellowstoneTransactionSignatureSelector::new("contains-0-O-I-l").is_err()); +} + +#[test] +fn yellowstone_transaction_update_decodes_current_storage_wire_including_v1_config_and_meta() { + let confirmed = yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionInfo { + signature: vec![9_u8; 64], + is_vote: false, + transaction: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::Transaction { + signatures: vec![vec![9_u8; 64], vec![8_u8; 64]], + message: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::Message { + header: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::MessageHeader { + num_required_signatures: 2, + num_readonly_signed_accounts: 1, + num_readonly_unsigned_accounts: 1, + }), + account_keys: vec![vec![1_u8; 32], vec![2_u8; 32]], + recent_blockhash: vec![3_u8; 32], + instructions: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::CompiledInstruction { + program_id_index: 1, + accounts: vec![0_u8, 1], + data: vec![4_u8, 5, 6], + }], + versioned: true, + address_table_lookups: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::MessageAddressTableLookup { + account_key: vec![4_u8; 32], + writable_indexes: vec![1_u8, 2], + readonly_indexes: vec![3_u8], + }], + config: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionConfig { + priority_fee: std::option::Option::Some(7), + compute_unit_limit: std::option::Option::Some(8), + loaded_accounts_data_size_limit: std::option::Option::Some(9), + heap_size: std::option::Option::Some(10), + }), + }), + }), + meta: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionStatusMeta { + err: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionError { err: vec![11_u8, 12] }), + fee: 5_000, + pre_balances: vec![100, 200], + post_balances: vec![90, 210], + inner_instructions: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::InnerInstructions { + index: 0, + instructions: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::InnerInstruction { + program_id_index: 1, + accounts: vec![0_u8], + data: vec![13_u8, 14], + stack_height: std::option::Option::Some(2), + }], + }], + inner_instructions_none: false, + log_messages: vec!["Program log: fixture".to_owned()], + log_messages_none: false, + pre_token_balances: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::TokenBalance { + account_index: 0, + mint: "mint-fixture".to_owned(), + ui_token_amount: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::UiTokenAmount { + ui_amount: 1.5, + decimals: 6, + amount: "1500000".to_owned(), + ui_amount_string: "1.5".to_owned(), + }), + owner: "owner-fixture".to_owned(), + program_id: "program-fixture".to_owned(), + }], + post_token_balances: vec![], + rewards: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::Reward { + pubkey: ksp_core_lib::Pubkey::new_from_array([5_u8; 32]).to_string(), + lamports: 17, + post_balance: 18, + reward_type: yellowstone_grpc_proto::solana::storage::confirmed_block::RewardType::Staking as i32, + commission: "5".to_owned(), + commission_bps: "500".to_owned(), + }], + loaded_writable_addresses: vec![vec![6_u8; 32]], + loaded_readonly_addresses: vec![vec![7_u8; 32]], + return_data: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::ReturnData { + program_id: vec![8_u8; 32], + data: vec![15_u8, 16], + }), + return_data_none: false, + compute_units_consumed: std::option::Option::Some(123), + cost_units: std::option::Option::Some(456), + }), + index: 3, + }; + let wire = yellowstone_grpc_proto::geyser::SubscribeUpdate { + filters: vec!["transactions-main".to_owned()], + update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::Transaction( + yellowstone_grpc_proto::geyser::SubscribeUpdateTransaction { transaction: std::option::Option::Some(confirmed), slot: 42 }, + )), + created_at: std::option::Option::Some(yellowstone_grpc_proto::prost_types::Timestamp { seconds: 100, nanos: 200 }), + }; + let update = super::decode_transaction_update(wire).expect("transaction update fixture must decode"); + assert_eq!(update.slot(), 42); + assert_eq!(update.filters()[0].as_str(), "transactions-main"); + assert_eq!(update.transaction().signature().as_bytes(), &[9_u8; 64]); + assert_eq!(format!("{:?}", update.transaction().signature()), "YellowstoneTransactionSignature()"); + assert_eq!(update.transaction().index(), 3); + assert_eq!(update.transaction().transaction().signatures().len(), 2); + assert!(update.transaction().transaction().message().versioned()); + let config = update.transaction().transaction().message().config().expect("v1 config must be preserved"); + assert_eq!(config.priority_fee(), std::option::Option::Some(7)); + assert_eq!(config.heap_size(), std::option::Option::Some(10)); + assert_eq!(update.transaction().meta().fee(), 5_000); + assert_eq!(update.transaction().meta().error().expect("error must be present").as_bytes(), &[11_u8, 12]); + assert_eq!(update.transaction().meta().inner_instructions()[0].instructions()[0].stack_height(), std::option::Option::Some(2)); + assert_eq!(update.transaction().meta().pre_token_balances()[0].ui_token_amount().expect("token amount must be present").amount(), "1500000"); + let token_debug = format!("{:?}", update.transaction().meta().pre_token_balances()[0]); + assert!(!token_debug.contains("mint-fixture")); + assert!(!token_debug.contains("owner-fixture")); + assert!(!token_debug.contains("program-fixture")); + assert_eq!(update.transaction().meta().rewards()[0].reward_type(), crate::YellowstoneRewardType::Staking); + assert_eq!(update.transaction().meta().loaded_writable_addresses()[0], ksp_core_lib::Pubkey::new_from_array([6_u8; 32])); + assert_eq!(update.transaction().meta().return_data().expect("return data must be present").data(), &[15_u8, 16]); + assert_eq!(update.transaction().meta().compute_units_consumed(), std::option::Option::Some(123)); + assert_eq!(update.transaction().meta().cost_units(), std::option::Option::Some(456)); + let debug = format!("{update:?}"); + assert!(!debug.contains("Program log: fixture")); + assert!(!debug.contains("11, 12")); +} + +#[test] +fn yellowstone_transaction_status_update_preserves_error_and_rejects_malformed_signature() { + let wire = yellowstone_grpc_proto::geyser::SubscribeUpdate { + filters: vec!["status".to_owned()], + update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::TransactionStatus( + yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionStatus { + slot: 55, + signature: vec![2_u8; 64], + is_vote: true, + index: 4, + err: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionError { err: vec![99_u8] }), + }, + )), + created_at: std::option::Option::None, + }; + let update = super::decode_transaction_status_update(wire).expect("transaction-status update must decode"); + assert_eq!(update.slot(), 55); + assert!(update.is_vote()); + assert_eq!(update.index(), 4); + assert_eq!(update.signature().as_bytes(), &[2_u8; 64]); + assert_eq!(update.error().expect("status error must be present").as_bytes(), &[99_u8]); + assert!(!format!("{update:?}").contains("99")); + let malformed = yellowstone_grpc_proto::geyser::SubscribeUpdate { + filters: vec!["status".to_owned()], + update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::TransactionStatus( + yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionStatus { + slot: 1, + signature: vec![0_u8; 63], + is_vote: false, + index: 0, + err: std::option::Option::None, + }, + )), + created_at: std::option::Option::None, + }; + assert!(super::decode_transaction_status_update(malformed).is_err()); +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/pool.rs b/crates/ksp-onchain-transport-lib/unit_tests/pool.rs deleted file mode 100644 index 135c1a5..0000000 --- a/crates/ksp-onchain-transport-lib/unit_tests/pool.rs +++ /dev/null @@ -1,336 +0,0 @@ -// file: crates/ksp-onchain-transport-lib/unit_tests/pool.rs -// version: 4 - -fn role(name: &str, priority: u32, request_kinds: std::vec::Vec) -> crate::HttpEndpointRoleSettings { - return crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new(name), - true, - request_kinds, - priority, - crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), - ); -} - -fn endpoint(name: &str, enabled: bool, priority: u32, request_kinds: std::vec::Vec) -> crate::HttpEndpointSettings { - return crate::HttpEndpointSettings::new( - name, - enabled, - crate::HttpProviderName::new("provider"), - crate::HttpClusterName::new("devnet"), - crate::HttpEndpointUrl::parse(format!("https://{name}.invalid/rpc?token=SECRET-CANARY")).expect("test URL must parse"), - std::time::Duration::from_secs(1), - std::time::Duration::from_secs(2), - std::option::Option::Some(4), - std::vec![role("default", priority, request_kinds)], - ); -} - -fn non_zero(value: u32) -> std::num::NonZeroU32 { - return std::num::NonZeroU32::new(value).expect("test limit must be non-zero"); -} - -fn limited_endpoint( - name: &str, - priority: u32, - requests_per_second: std::option::Option, - burst_capacity: std::option::Option, - max_concurrent_requests: std::option::Option, - cooldown: std::option::Option, -) -> crate::HttpEndpointSettings { - let limits = crate::HttpRoleLimits::new( - requests_per_second.map(|value| return non_zero(value)), - burst_capacity.map(|value| return non_zero(value)), - max_concurrent_requests.map(|value| return non_zero(value)), - cooldown, - ); - let role = crate::HttpEndpointRoleSettings::new(crate::HttpRoleName::new("default"), true, std::vec![crate::HttpRequestKind::wildcard()], priority, limits); - return crate::HttpEndpointSettings::new( - name, - true, - crate::HttpProviderName::new("provider"), - crate::HttpClusterName::new("devnet"), - crate::HttpEndpointUrl::parse(format!("https://{name}.invalid/rpc?token=SECRET-CANARY")).expect("test URL must parse"), - std::time::Duration::from_secs(1), - std::time::Duration::from_secs(2), - std::option::Option::Some(4), - std::vec![role], - ); -} - -fn settings(endpoints: std::vec::Vec) -> crate::HttpTransportSettings { - return crate::HttpTransportSettings::new( - endpoints, - crate::HttpRetrySettings::new(2, std::time::Duration::from_millis(10), std::time::Duration::from_millis(50)), - ); -} - -#[test] -fn pool_prefers_lowest_priority_tier() { - let pool = crate::HttpTransportPool::new(settings(std::vec![ - endpoint("secondary", true, 20, std::vec![crate::HttpRequestKind::wildcard()]), - endpoint("primary", true, 10, std::vec![crate::HttpRequestKind::wildcard()]), - ])) - .expect("pool must build"); - let selection = pool - .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")) - .expect("selection must succeed"); - assert_eq!(selection.endpoint_name(), "primary"); - assert_eq!(selection.priority(), 10); -} - -#[test] -fn pool_round_robins_fairly_inside_best_priority_tier() { - let pool = crate::HttpTransportPool::new(settings(std::vec![ - endpoint("one", true, 10, std::vec![crate::HttpRequestKind::wildcard()]), - endpoint("two", true, 10, std::vec![crate::HttpRequestKind::wildcard()]), - endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]), - ])) - .expect("pool must build"); - let role = crate::HttpRoleName::new("default"); - let kind = crate::HttpRequestKind::new("get_balance"); - let first = pool.select_for_request_kind(&role, &kind).expect("first selection must succeed"); - let second = pool.select_for_request_kind(&role, &kind).expect("second selection must succeed"); - let third = pool.select_for_request_kind(&role, &kind).expect("third selection must succeed"); - assert_eq!(first.endpoint_name(), "one"); - assert_eq!(second.endpoint_name(), "two"); - assert_eq!(third.endpoint_name(), "one"); -} - -#[test] -fn disabled_best_priority_endpoint_falls_back_to_next_tier() { - let pool = crate::HttpTransportPool::new(settings(std::vec![ - endpoint("disabled-primary", false, 1, std::vec![crate::HttpRequestKind::wildcard()]), - endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]), - ])) - .expect("pool must build"); - let selection = pool - .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")) - .expect("fallback must be selected"); - assert_eq!(selection.endpoint_name(), "fallback"); -} - -#[test] -fn pool_filters_role_and_capability_before_priority() { - let pool = crate::HttpTransportPool::new(settings(std::vec![ - endpoint("wrong-capability", true, 1, std::vec![crate::HttpRequestKind::new("send_transaction")]), - endpoint("matching", true, 50, std::vec![crate::HttpRequestKind::new("get_balance")]), - ])) - .expect("pool must build"); - let selection = pool - .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")) - .expect("matching capability must be selected"); - assert_eq!(selection.endpoint_name(), "matching"); -} - -#[test] -fn pool_returns_structured_error_when_no_endpoint_matches() { - let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("read-only", true, 10, std::vec![crate::HttpRequestKind::new("get_balance")],)])) - .expect("pool must build"); - let error = pool - .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("send_transaction")) - .expect_err("unsupported request kind must fail selection"); - assert_eq!(error.code(), crate::ERROR_CODE_ENDPOINT_SELECTION_FAILED); - assert!(!format!("{error:?}").contains("SECRET-CANARY")); -} - -#[test] -fn standard_method_selection_uses_registry_request_kind() { - let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("balance", true, 10, std::vec![crate::HttpRequestKind::new("get_balance")],)])) - .expect("pool must build"); - let method = crate::find_http_rpc_method("getBalance").expect("audited method must exist"); - let selection = pool.select_for_method(&crate::HttpRoleName::new("default"), method).expect("standard method must route"); - assert_eq!(selection.request_kind().as_str(), "get_balance"); -} - -#[test] -fn pool_snapshot_is_safe_and_preserves_disabled_endpoints() { - let pool = crate::HttpTransportPool::new(settings(std::vec![ - endpoint("enabled", true, 10, std::vec![crate::HttpRequestKind::wildcard()]), - endpoint("disabled", false, 10, std::vec![crate::HttpRequestKind::wildcard()]), - ])) - .expect("pool must build"); - let snapshot = pool.snapshot(); - let rendered = format!("{snapshot:?} {pool:?}"); - assert_eq!(snapshot.endpoint_count(), 2); - assert_eq!(snapshot.available_endpoint_count(), 1); - assert!(!rendered.contains("SECRET-CANARY")); - assert!(!rendered.contains(".invalid/rpc")); -} - -#[test] -fn disabled_role_is_excluded_before_priority_selection() { - let base = endpoint("disabled-role", true, 1, std::vec![crate::HttpRequestKind::wildcard()]); - let disabled_role = crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new("default"), - false, - std::vec![crate::HttpRequestKind::wildcard()], - 1, - crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), - ); - let enabled_non_matching_role = crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new("maintenance"), - true, - std::vec![crate::HttpRequestKind::wildcard()], - 1, - crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), - ); - let disabled_role_endpoint = crate::HttpEndpointSettings::new( - base.name(), - true, - base.provider().clone(), - base.cluster().clone(), - base.url().clone(), - base.connect_timeout(), - base.request_timeout(), - base.max_idle_connections_per_host(), - std::vec![disabled_role, enabled_non_matching_role], - ); - let pool = crate::HttpTransportPool::new(settings(std::vec![ - disabled_role_endpoint, - endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]), - ])) - .expect("pool must build"); - let selection = pool - .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")) - .expect("enabled fallback role must be selected"); - assert_eq!(selection.endpoint_name(), "fallback"); -} - -#[test] -fn removed_standard_method_is_rejected_before_endpoint_routing() { - let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("wildcard", true, 10, std::vec![crate::HttpRequestKind::wildcard()],)])) - .expect("pool must build"); - let method = crate::find_http_rpc_method("confirmTransaction").expect("historical method must exist"); - let error = pool.select_for_method(&crate::HttpRoleName::new("default"), method).expect_err("removed standard method must be rejected before routing"); - assert_eq!(error.code(), crate::ERROR_CODE_METHOD_REMOVED); -} - -#[tokio::test] -async fn runtime_concurrency_saturation_falls_back_to_lower_priority_tier() { - let pool = crate::HttpTransportPool::new(settings(std::vec![ - limited_endpoint("primary", 1, std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None), - limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None), - ])) - .expect("pool must build"); - let role = crate::HttpRoleName::new("default"); - let kind = crate::HttpRequestKind::new("get_balance"); - let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first request must acquire primary"); - assert_eq!(first.selection().endpoint_name(), "primary"); - let second = pool.acquire_for_request_kind(&role, &kind).await.expect("second request must fall back while primary is saturated"); - assert_eq!(second.selection().endpoint_name(), "fallback"); -} - -#[tokio::test] -async fn runtime_token_bucket_exhaustion_falls_back_without_busy_waiting() { - let pool = crate::HttpTransportPool::new(settings(std::vec![ - limited_endpoint("primary", 1, std::option::Option::Some(1), std::option::Option::Some(1), std::option::Option::None, std::option::Option::None), - limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), - ])) - .expect("pool must build"); - let role = crate::HttpRoleName::new("default"); - let kind = crate::HttpRequestKind::new("get_balance"); - let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first request must consume primary token"); - assert_eq!(first.selection().endpoint_name(), "primary"); - drop(first); - let second = pool.acquire_for_request_kind(&role, &kind).await.expect("fallback must be used while primary token bucket refills"); - assert_eq!(second.selection().endpoint_name(), "fallback"); -} - -#[tokio::test] -async fn provider_cooldown_excludes_rate_limited_role_and_uses_fallback() { - let pool = crate::HttpTransportPool::new(settings(std::vec![ - limited_endpoint( - "primary", - 1, - std::option::Option::None, - std::option::Option::None, - std::option::Option::None, - std::option::Option::Some(std::time::Duration::from_millis(50)), - ), - limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), - ])) - .expect("pool must build"); - let role = crate::HttpRoleName::new("default"); - let kind = crate::HttpRequestKind::new("get_balance"); - let primary = pool.acquire_for_request_kind(&role, &kind).await.expect("primary must be acquired"); - assert_eq!(primary.selection().endpoint_name(), "primary"); - let pause = primary.record_rate_limited(std::option::Option::None); - assert_eq!(pause, std::time::Duration::from_millis(50)); - drop(primary); - let fallback = pool.acquire_for_request_kind(&role, &kind).await.expect("fallback must be selected during primary cooldown"); - assert_eq!(fallback.selection().endpoint_name(), "fallback"); -} - -#[tokio::test] -async fn admission_waits_for_released_concurrency_without_holding_a_sync_mutex_across_await() { - let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint( - "primary", - 1, - std::option::Option::None, - std::option::Option::None, - std::option::Option::Some(1), - std::option::Option::None, - )])) - .expect("pool must build"); - let role = crate::HttpRoleName::new("default"); - let kind = crate::HttpRequestKind::new("get_balance"); - let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first permit must be acquired"); - let release_task = tokio::spawn(async move { - tokio::time::sleep(std::time::Duration::from_millis(10)).await; - drop(first); - }); - let second = pool - .acquire_for_request_kind_with_timeout(&role, &kind, std::time::Duration::from_millis(100)) - .await - .expect("second permit must wake after concurrency release"); - assert_eq!(second.selection().endpoint_name(), "primary"); - release_task.await.expect("release task must complete"); -} - -#[tokio::test] -async fn admission_timeout_is_bounded_when_concurrency_never_becomes_available() { - let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint( - "primary", - 1, - std::option::Option::None, - std::option::Option::None, - std::option::Option::Some(1), - std::option::Option::None, - )])) - .expect("pool must build"); - let role = crate::HttpRoleName::new("default"); - let kind = crate::HttpRequestKind::new("get_balance"); - let _held = pool.acquire_for_request_kind(&role, &kind).await.expect("first permit must be acquired"); - let error = pool - .acquire_for_request_kind_with_timeout(&role, &kind, std::time::Duration::from_millis(20)) - .await - .expect_err("second permit must time out while concurrency remains saturated"); - assert_eq!(error.code(), crate::ERROR_CODE_TIMEOUT); -} - -#[tokio::test] -async fn passive_health_snapshot_moves_from_degraded_back_to_available_after_success() { - let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint( - "primary", - 1, - std::option::Option::None, - std::option::Option::None, - std::option::Option::None, - std::option::Option::None, - )])) - .expect("pool must build"); - let role = crate::HttpRoleName::new("default"); - let kind = crate::HttpRequestKind::new("get_balance"); - let first = pool.acquire_for_request_kind(&role, &kind).await.expect("request permit must be acquired"); - first.record_failure(); - drop(first); - let degraded = pool.snapshot(); - assert_eq!(degraded.endpoints()[0].availability(), crate::HttpEndpointAvailability::Degraded); - assert_eq!(degraded.endpoints()[0].roles()[0].failure_count(), 1); - let second = pool.acquire_for_request_kind(&role, &kind).await.expect("degraded endpoint remains eligible for passive recovery"); - second.record_success(); - drop(second); - let recovered = pool.snapshot(); - assert_eq!(recovered.endpoints()[0].availability(), crate::HttpEndpointAvailability::Available); - assert_eq!(recovered.endpoints()[0].roles()[0].success_count(), 1); -} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/resilience.rs b/crates/ksp-onchain-transport-lib/unit_tests/resilience.rs deleted file mode 100644 index dc82a3c..0000000 --- a/crates/ksp-onchain-transport-lib/unit_tests/resilience.rs +++ /dev/null @@ -1,186 +0,0 @@ -// file: crates/ksp-onchain-transport-lib/unit_tests/resilience.rs -// version: 2 - -fn non_zero(value: u32) -> std::num::NonZeroU32 { - return std::num::NonZeroU32::new(value).expect("test limit must be non-zero"); -} - -fn retry_settings() -> crate::HttpRetrySettings { - return crate::HttpRetrySettings::new(4, std::time::Duration::from_millis(100), std::time::Duration::from_millis(500)); -} - -fn method(name: &str) -> &'static crate::HttpRpcMethodDescriptor { - return crate::find_http_rpc_method(name).expect("audited test method must exist"); -} - -fn role_limits( - requests_per_second: std::option::Option, - burst_capacity: std::option::Option, - max_concurrent_requests: std::option::Option, - cooldown: std::option::Option, -) -> crate::HttpEndpointRoleSettings { - return crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new("default"), - true, - std::vec![crate::HttpRequestKind::wildcard()], - 10, - crate::HttpRoleLimits::new( - requests_per_second.map(|value| return non_zero(value)), - burst_capacity.map(|value| return non_zero(value)), - max_concurrent_requests.map(|value| return non_zero(value)), - cooldown, - ), - ); -} - -#[test] -fn retry_backoff_is_exponential_and_bounded() { - let settings = retry_settings(); - assert_eq!(super::retry_backoff(&settings, 1), std::time::Duration::from_millis(100)); - assert_eq!(super::retry_backoff(&settings, 2), std::time::Duration::from_millis(200)); - assert_eq!(super::retry_backoff(&settings, 3), std::time::Duration::from_millis(400)); - assert_eq!(super::retry_backoff(&settings, 4), std::time::Duration::from_millis(500)); - assert_eq!(super::retry_backoff(&settings, 32), std::time::Duration::from_millis(500)); -} - -#[test] -fn retry_safe_timeout_is_retried_until_budget_is_exhausted() { - let settings = retry_settings(); - let first = crate::evaluate_transport_retry( - method("getBalance"), - &settings, - crate::HttpRetryCause::Timeout, - crate::HttpDispatchState::DispatchedAmbiguous, - 0, - std::option::Option::None, - ); - assert_eq!(first, crate::HttpRetryDecision::RetryAfter(std::time::Duration::from_millis(100))); - let exhausted = crate::evaluate_transport_retry( - method("getBalance"), - &settings, - crate::HttpRetryCause::Timeout, - crate::HttpDispatchState::DispatchedAmbiguous, - settings.max_retries(), - std::option::Option::None, - ); - assert_eq!(exhausted, crate::HttpRetryDecision::Stop); -} - -#[test] -fn write_submission_never_retries_after_ambiguous_dispatch() { - let decision = crate::evaluate_transport_retry( - method("sendTransaction"), - &retry_settings(), - crate::HttpRetryCause::Connection, - crate::HttpDispatchState::DispatchedAmbiguous, - 0, - std::option::Option::None, - ); - assert_eq!(decision, crate::HttpRetryDecision::Stop); -} - -#[test] -fn write_submission_can_retry_when_transport_proves_no_dispatch() { - let decision = crate::evaluate_transport_retry( - method("sendTransaction"), - &retry_settings(), - crate::HttpRetryCause::Connection, - crate::HttpDispatchState::NotDispatched, - 0, - std::option::Option::None, - ); - assert!(decision.should_retry()); -} - -#[test] -fn rpc_application_and_invalid_response_are_not_transport_retries() { - for cause in [crate::HttpRetryCause::RpcApplication, crate::HttpRetryCause::InvalidResponse, crate::HttpRetryCause::Request] { - let decision = crate::evaluate_transport_retry( - method("getBalance"), - &retry_settings(), - cause, - crate::HttpDispatchState::NotDispatched, - 0, - std::option::Option::None, - ); - assert_eq!(decision, crate::HttpRetryDecision::Stop); - } -} - -#[test] -fn provider_retry_after_can_extend_backoff_but_is_defensively_bounded() { - let settings = retry_settings(); - let extended = crate::evaluate_transport_retry( - method("getBalance"), - &settings, - crate::HttpRetryCause::RateLimited, - crate::HttpDispatchState::DispatchedAmbiguous, - 0, - std::option::Option::Some(std::time::Duration::from_secs(3)), - ); - assert_eq!(extended.delay(), std::option::Option::Some(std::time::Duration::from_secs(3))); - let bounded = crate::evaluate_transport_retry( - method("getBalance"), - &settings, - crate::HttpRetryCause::RateLimited, - crate::HttpDispatchState::DispatchedAmbiguous, - 0, - std::option::Option::Some(std::time::Duration::from_secs(600)), - ); - assert_eq!(bounded.delay(), std::option::Option::Some(std::time::Duration::from_secs(60))); -} - -#[test] -fn token_bucket_consumes_burst_then_refills_from_elapsed_time() { - let start = std::time::Instant::now(); - let mut bucket = super::HttpTokenBucketState::new(2, 2, start); - assert!(bucket.try_consume_at(start).is_none()); - assert!(bucket.try_consume_at(start).is_none()); - assert!(bucket.try_consume_at(start).is_some()); - let later = start.checked_add(std::time::Duration::from_millis(500)).expect("test instant must advance"); - assert!(bucket.try_consume_at(later).is_none()); -} - -#[test] -fn absent_burst_capacity_defaults_to_one_second_of_rps_capacity() { - let role = role_limits(std::option::Option::Some(2), std::option::Option::None, std::option::Option::None, std::option::Option::None); - let runtime = std::sync::Arc::new(crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new()))); - let now = std::time::Instant::now(); - let first = runtime.try_acquire(now); - let second = runtime.try_acquire(now); - let third = runtime.try_acquire(now); - assert!(matches!(first, crate::RoleAdmissionAttempt::Ready(_))); - assert!(matches!(second, crate::RoleAdmissionAttempt::Ready(_))); - assert!(matches!(third, crate::RoleAdmissionAttempt::BlockedUntil(_))); -} - -#[test] -fn concurrency_semaphore_releases_capacity_when_permit_is_dropped() { - let role = role_limits(std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None); - let runtime = std::sync::Arc::new(crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new()))); - let now = std::time::Instant::now(); - let first = runtime.try_acquire(now); - let held = match first { - crate::RoleAdmissionAttempt::Ready(permit) => permit, - _ => panic!("first concurrency permit must be available"), - }; - assert!(matches!(runtime.try_acquire(now), crate::RoleAdmissionAttempt::ConcurrencySaturated)); - drop(held); - assert!(matches!(runtime.try_acquire(now), crate::RoleAdmissionAttempt::Ready(_))); -} - -#[test] -fn rate_limit_cooldown_marks_role_and_caps_provider_delay() { - let role = role_limits( - std::option::Option::None, - std::option::Option::None, - std::option::Option::None, - std::option::Option::Some(std::time::Duration::from_millis(10)), - ); - let runtime = crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new())); - let pause = runtime.record_rate_limited(std::option::Option::Some(std::time::Duration::from_secs(600))); - assert_eq!(pause, std::time::Duration::from_secs(60)); - assert_eq!(runtime.rate_limit_count(), 1); - assert_eq!(runtime.failure_count(), 1); - assert_eq!(runtime.availability(std::time::Instant::now()), crate::HttpEndpointAvailability::RateLimited); -} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/settings.rs b/crates/ksp-onchain-transport-lib/unit_tests/settings.rs deleted file mode 100644 index be1d37a..0000000 --- a/crates/ksp-onchain-transport-lib/unit_tests/settings.rs +++ /dev/null @@ -1,207 +0,0 @@ -// file: crates/ksp-onchain-transport-lib/unit_tests/settings.rs -// version: 2 - -fn non_zero(value: u32) -> std::num::NonZeroU32 { - return std::num::NonZeroU32::new(value).expect("test non-zero value must remain non-zero"); -} - -fn valid_settings(url_text: &str) -> crate::HttpTransportSettings { - let url = crate::HttpEndpointUrl::parse(url_text).expect("test URL must be valid"); - let limits = crate::HttpRoleLimits::new( - std::option::Option::Some(non_zero(10)), - std::option::Option::Some(non_zero(20)), - std::option::Option::Some(non_zero(4)), - std::option::Option::Some(std::time::Duration::from_millis(500)), - ); - let role = crate::HttpEndpointRoleSettings::new(crate::HttpRoleName::new("default"), true, std::vec![crate::HttpRequestKind::wildcard()], 100, limits); - let endpoint = crate::HttpEndpointSettings::new( - "devnet_public", - true, - crate::HttpProviderName::new("solana-public"), - crate::HttpClusterName::new("devnet"), - url, - std::time::Duration::from_secs(5), - std::time::Duration::from_secs(15), - std::option::Option::Some(8), - std::vec![role], - ); - return crate::HttpTransportSettings::new( - std::vec![endpoint], - crate::HttpRetrySettings::new(2, std::time::Duration::from_millis(100), std::time::Duration::from_secs(2)), - ); -} - -#[test] -fn endpoint_url_accepts_http_and_https() { - assert!(crate::HttpEndpointUrl::parse("https://api.devnet.solana.com").is_ok()); - assert!(crate::HttpEndpointUrl::parse("http://127.0.0.1:8899").is_ok()); -} - -#[test] -fn endpoint_url_rejects_non_http_schemes() { - let result = crate::HttpEndpointUrl::parse("ws://api.devnet.solana.com"); - let error = result.expect_err("WebSocket URL must not be accepted by HTTP settings"); - assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS); -} - -#[test] -fn endpoint_url_debug_redacts_secret_material() { - let url = crate::HttpEndpointUrl::parse("https://provider.invalid/rpc?api-key=SECRET-CANARY").expect("test URL must parse"); - let rendered = format!("{url:?}"); - assert!(rendered.contains("")); - assert!(!rendered.contains("SECRET-CANARY")); - assert!(!rendered.contains("provider.invalid")); -} - -#[test] -fn valid_transport_settings_pass_validation() { - let settings = valid_settings("https://api.devnet.solana.com"); - assert!(settings.validate().is_ok()); -} - -#[test] -fn transport_settings_debug_does_not_leak_endpoint_url() { - let settings = valid_settings("https://provider.invalid/rpc?api-key=SECRET-CANARY"); - let rendered = format!("{settings:?}"); - assert!(!rendered.contains("SECRET-CANARY")); - assert!(!rendered.contains("provider.invalid")); - assert!(rendered.contains("HttpEndpointUrl()")); -} - -#[test] -fn transport_settings_require_one_enabled_endpoint() { - let url = crate::HttpEndpointUrl::parse("https://api.devnet.solana.com").expect("test URL must parse"); - let role = crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new("default"), - true, - std::vec![crate::HttpRequestKind::wildcard()], - 100, - crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), - ); - let endpoint = crate::HttpEndpointSettings::new( - "disabled", - false, - crate::HttpProviderName::new("provider"), - crate::HttpClusterName::new("devnet"), - url, - std::time::Duration::from_secs(1), - std::time::Duration::from_secs(1), - std::option::Option::None, - std::vec![role], - ); - let settings = crate::HttpTransportSettings::new( - std::vec![endpoint], - crate::HttpRetrySettings::new(1, std::time::Duration::from_millis(1), std::time::Duration::from_millis(2)), - ); - let error = settings.validate().expect_err("all-disabled settings must fail"); - assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS); -} - -#[test] -fn transport_settings_reject_duplicate_endpoint_names() { - let first = valid_settings("https://one.invalid"); - let second = valid_settings("https://two.invalid"); - let settings = crate::HttpTransportSettings::new(std::vec![first.endpoints()[0].clone(), second.endpoints()[0].clone()], first.retry().clone()); - let error = settings.validate().expect_err("duplicate endpoint names must fail"); - assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS); -} - -#[test] -fn transport_settings_reject_duplicate_roles() { - let base = valid_settings("https://api.devnet.solana.com"); - let endpoint = &base.endpoints()[0]; - let duplicated_endpoint = crate::HttpEndpointSettings::new( - endpoint.name(), - true, - endpoint.provider().clone(), - endpoint.cluster().clone(), - endpoint.url().clone(), - endpoint.connect_timeout(), - endpoint.request_timeout(), - endpoint.max_idle_connections_per_host(), - std::vec![endpoint.roles()[0].clone(), endpoint.roles()[0].clone()], - ); - let settings = crate::HttpTransportSettings::new(std::vec![duplicated_endpoint], base.retry().clone()); - assert!(settings.validate().is_err()); -} - -#[test] -fn transport_settings_reject_wildcard_mixed_with_specific_kind() { - let base = valid_settings("https://api.devnet.solana.com"); - let endpoint = &base.endpoints()[0]; - let role = crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new("default"), - true, - std::vec![crate::HttpRequestKind::wildcard(), crate::HttpRequestKind::new("get_balance")], - 100, - endpoint.roles()[0].limits().clone(), - ); - let modified_endpoint = crate::HttpEndpointSettings::new( - endpoint.name(), - true, - endpoint.provider().clone(), - endpoint.cluster().clone(), - endpoint.url().clone(), - endpoint.connect_timeout(), - endpoint.request_timeout(), - endpoint.max_idle_connections_per_host(), - std::vec![role], - ); - let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone()); - assert!(settings.validate().is_err()); -} - -#[test] -fn transport_settings_reject_burst_without_rps() { - let base = valid_settings("https://api.devnet.solana.com"); - let endpoint = &base.endpoints()[0]; - let role = crate::HttpEndpointRoleSettings::new( - crate::HttpRoleName::new("default"), - true, - std::vec![crate::HttpRequestKind::wildcard()], - 100, - crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::Some(non_zero(2)), std::option::Option::None, std::option::Option::None), - ); - let modified_endpoint = crate::HttpEndpointSettings::new( - endpoint.name(), - true, - endpoint.provider().clone(), - endpoint.cluster().clone(), - endpoint.url().clone(), - endpoint.connect_timeout(), - endpoint.request_timeout(), - endpoint.max_idle_connections_per_host(), - std::vec![role], - ); - let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone()); - assert!(settings.validate().is_err()); -} - -#[test] -fn transport_settings_reject_reversed_retry_backoff() { - let base = valid_settings("https://api.devnet.solana.com"); - let settings = crate::HttpTransportSettings::new( - base.endpoints().to_vec(), - crate::HttpRetrySettings::new(2, std::time::Duration::from_secs(2), std::time::Duration::from_secs(1)), - ); - assert!(settings.validate().is_err()); -} - -#[test] -fn transport_settings_reject_zero_request_timeout() { - let base = valid_settings("https://api.devnet.solana.com"); - let endpoint = &base.endpoints()[0]; - let modified_endpoint = crate::HttpEndpointSettings::new( - endpoint.name(), - true, - endpoint.provider().clone(), - endpoint.cluster().clone(), - endpoint.url().clone(), - endpoint.connect_timeout(), - std::time::Duration::ZERO, - endpoint.max_idle_connections_per_host(), - endpoint.roles().to_vec(), - ); - let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone()); - assert!(settings.validate().is_err()); -} diff --git a/deltas/0.2.9/pre.007.md b/deltas/0.2.9/pre.007.md new file mode 100644 index 0000000..aaef276 --- /dev/null +++ b/deltas/0.2.9/pre.007.md @@ -0,0 +1,165 @@ + + + +# Delta `0.2.9-pre.007` — Yellowstone Transactions + `transaction_status` + +## 1. Base et gate précédent + +Base exacte : + +```text +0.2.9-pre.006 +``` + +Le gate opérateur `pre.006` est fermé : fmt/audit/check/Clippy/workspace PASS, Transport 364 unit + 45 public API + 39 release-completeness + 4 doctests, dependency canary Core 3/3 PASS. + +Version technique de cette tranche : + +```text +0.2.9-pre.7 +``` + +Commit attendu après gate vert : + +```text +v0.2.9-pre.007 +``` + +Aucun tag Git stable n'est créé pour cette prerelease. + +## 2. Objet de la tranche + +`pre.007` complète exclusivement les deux familles Yellowstone standard `transactions` et `transactions_status` du `SubscribeRequest`/`SubscribeUpdate` courant. Elle n'ouvre toujours pas le stream bidi runtime. + +Le proto publié `yellowstone-grpc-proto 12.6.0` a été réaudité avant implémentation. Le même `SubscribeRequestFilterTransactions` est utilisé par les deux maps et expose : + +```text +vote? +failed? +signature? +account_include[] +account_exclude[] +account_required[] +cuckoo_account_include? +token_accounts? = ALL | BALANCE_CHANGED +``` + +## 3. Contrat request KSP + +La tranche matérialise `YellowstoneSubscribeTransactionFilter` et les types associés : + +```text +YellowstoneTransactionSignatureSelector +YellowstoneTokenAccountExpansion +YellowstoneCuckooFilter réutilisé +``` + +La signature textuelle est : + +- non vide ; +- bornée à 128 octets ; +- strictement Base58 ; +- décodée par KSP pour vérifier une largeur exacte de 64 octets ; +- absente des diagnostics et du `Debug`. + +Les listes include/exclude/required conservent l'ordre d'insertion et réutilisent le bound déterministe des sélecteurs Accounts. Le Cuckoo transaction réutilise le contrat standard KSP déjà introduit en `pre.005`. + +## 4. Updates et `solana-storage.proto` + +Les nouveaux DTOs KSP couvrent les variantes : + +```text +SubscribeUpdateTransaction +SubscribeUpdateTransactionStatus +``` + +et projettent sans raw reexport upstream les structures utiles : + +```text +Transaction / Message / MessageHeader +CompiledInstruction / MessageAddressTableLookup +TransactionConfig V1 +TransactionStatusMeta / TransactionError +InnerInstructions / InnerInstruction +TokenBalance / UiTokenAmount +ReturnData / Reward +``` + +Le `Message.config` optionnel actuel est conservé avec : + +```text +priority_fee? +compute_unit_limit? +loaded_accounts_data_size_limit? +heap_size? +``` + +`TransactionStatusMeta` conserve séparément les marqueurs legacy `inner_instructions_none`, `log_messages_none` et `return_data_none`, ainsi que les payloads correspondants. Les champs optionnels `compute_units_consumed` et `cost_units`, les loaded addresses, rewards et token balances sont également préservés. + +`TransactionError.err` reste un payload opaque borné ; KSP n'introduit aucun décodage JSON ou Program arbitraire. + +## 5. Validation et sûreté + +Les décodeurs test-only vérifient notamment : + +| Élément | Politique `pre.007` | +|------------------------------------------------------|----------------------------------------------------------------------------| +| signatures wire | exactement 64 octets | +| hash / account keys / loaded addresses / program ids | exactement 32 octets | +| vecteurs transaction/meta | bornes KSP avant projection | +| instruction / return-data payloads | bornés | +| logs et textes provider | count/length bornés | +| error bytes | opaques et bornés | +| `Debug` | pas de signature, pubkey, log, instruction data ou error bytes arbitraires | + +Les conversions protobuf request et les décodeurs update restent `#[cfg(test)]` jusqu'à `pre.009`, premier consommateur runtime prévu lors de l'ouverture du stream bidi. + +## 6. Frontières préservées + +Toujours hors `pre.007` : + +```text +Blocks + block_meta + entry pre.008 +bidi/backpressure/Ping-Pong/shutdown pre.009 +reconnect/replay/gaps/duplicates pre.010 +Config V3 + PublicNode profiles pre.011 +PublicNode live/compliance pre.012 +SubscribeDeshred OUT 0.2.9 standard +``` + +Aucun provider N3 n'est ajouté. Aucune dépendance ni feature Cargo n'est modifiée. + +## 7. Preuves ajoutées + +Tests Transport ajoutés : + +```text +transaction filter exact wire + redaction + Base58 width +transaction update storage fixture + TransactionConfig V1 + meta +transaction_status update + error + malformed signature +public API root contract +release-completeness exact scope +``` + +Compte cible après compilation : + +```text +Transport unit ~367 +Transport public API 46 +Transport release completeness 40 +Transport doctests 4 +``` + +## 8. Gate opérateur requis + +```bash +cargo fmt --all +python3 scripts/audit_rust_workspace_rules.py +cargo check --workspace +cargo clippy --workspace --all-targets +cargo test -p ksp-onchain-transport-lib +cargo test -p ksp-core-lib --test workspace_dependencies +cargo test --workspace +``` + +Aucun `cargo tree` supplémentaire n'est requis : aucune dépendance ou feature ne change dans cette tranche. diff --git a/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md b/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md index c884d7e..f85ea4f 100644 --- a/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md +++ b/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md @@ -3,7 +3,7 @@ # Plan `0.2.9` — moteur Yellowstone gRPC + standard Solana + PublicNode -> **Statut : `0.2.9-pre.005-fix.001` est fermée sur gate opérateur intégralement vert et sans warning : fmt/audit/check/Clippy/workspace PASS, Transport 364 unit + 45 public API + 38 completeness + 4 doctests. `0.2.9-pre.006` est une tranche structurelle dédiée au nommage des modules privés : les cinq modules sans ambiguïté HTTP (`client`, `executor`, `pool`, `resilience`, `settings`) deviennent `http_*`, avec leurs unit tests miroirs. Les modules `rpc_*`, `json_rpc`, `constants` et `error` ne sont pas préfixés artificiellement car ils décrivent une couche/protocole partagé ou portent déjà des DTOs utilisés par WS/gRPC. Cette insertion décale l’ancien forecast fonctionnel `pre.006–011` vers `pre.007–012` sans changer son contenu. Aucun contrat public ni comportement runtime ne change.** +> **Statut : `0.2.9-pre.006` est fermée sur gate opérateur intégralement vert : fmt/audit/check/Clippy/workspace PASS, Transport 364 unit + 45 public API + 39 completeness + 4 doctests. `0.2.9-pre.007` est candidate et matérialise uniquement Transactions + `transaction_status` du `Subscribe` standard : filtres complets, signature Base58 décodée exactement sur 64 octets, Cuckoo/token-account expansion, DTOs transaction/meta complets et Transaction V1 `TransactionConfig`. Blocks restent `pre.008`; aucun stream bidi n’est encore ouvert.** ## 1. Objet, base et état d'ouverture @@ -925,11 +925,11 @@ pre.004 DONE — standard Solana : Subscribe foundation + maps/commitment/ping/ pre.005 DONE — standard Solana : Accounts + Slots filters/updates budget : 15–20 min ; gate final fix.001 : fmt/audit/check/Clippy/workspace PASS sans warning + Transport 364/45/38/4 -pre.006 CANDIDATE — structure Transport : namespace privé HTTP explicite - budget : 15–20 min ; preuve : 5 modules + 5 unit tests renommés `http_*`, API publique inchangée, canari d'absence des anciens modules +pre.006 DONE — structure Transport : namespace privé HTTP explicite + budget : 15–20 min ; gate opérateur final PASS 364/45/39/4 + workspace -pre.007 standard Solana : Transactions + transaction_status - budget : 15–20 min ; preuve : include/exclude/required/Cuckoo/token expansion + tx/meta +pre.007 CANDIDATE — standard Solana : Transactions + transaction_status + budget : 15–20 min ; preuve : include/exclude/required/Cuckoo/token expansion + tx/meta + TransactionConfig V1 pre.008 standard Solana : Blocks + block_meta + entry budget : 15–20 min ; preuve : counts/arrays/optional/oneof/payload bounds @@ -1216,3 +1216,60 @@ Cette tranche est volontairement séparée de Transactions afin de respecter le **Gate candidat :** audit Rust clean, public API inchangée, cinq anciens modules absents après suppression opérateur, cinq nouveaux modules `http_*` compilés, release-completeness canary dédié, workspace complet vert. + +## 21. Gate final `pre.006` — namespace privé HTTP + +Preuve opérateur du `2026-08-24` : + +| Gate | Résultat final | +|------------------------------------------|----------------| +| `cargo fmt --all` | PASS | +| audit Rust workspace | PASS / clean | +| `cargo check --workspace` | PASS | +| `cargo clippy --workspace --all-targets` | PASS | +| Transport unit | 364/364 PASS | +| Transport `public_api` | 45/45 PASS | +| Transport `release_completeness` | 39/39 PASS | +| Transport doctests | 4/4 PASS | +| Core dependency canary | 3/3 PASS | +| `cargo test --workspace` | PASS | + +Les cinq modules privés HTTP et leurs unit tests miroirs sont désormais effectivement nommés `http_*`. Le canari `release_v0_2_9_pre_006_namespaces_unambiguously_http_owned_private_modules` passe et les modules partagés `rpc_*`/`json_rpc` restent volontairement non préfixés. + +**Verdict : `pre.006` fermée.** + +## 22. `pre.007` — Transactions + `transaction_status` candidate + +Le proto `yellowstone-grpc-proto 12.6.0` est réaudité avant implémentation. `SubscribeRequestFilterTransactions` contient exactement `vote?`, `failed?`, `signature?`, `account_include[]`, `account_exclude[]`, `account_required[]`, `cuckoo_account_include?` et `token_accounts?`. La même structure filtre `transactions` et `transactions_status`. + +La candidate matérialise : + +```text +request filters + vote? / failed? + signature? : Base58 validé et décodé exactement sur 64 octets + account_include[] / account_exclude[] / account_required[] + cuckoo_account_include? + token_accounts? = ALL | BALANCE_CHANGED + +updates + transaction : filters + created_at + slot + signature/is_vote/index + transaction + meta + transaction_status : filters + created_at + slot + signature/is_vote/index + error? + +solana-storage typed + Transaction / Message / MessageHeader + CompiledInstruction / MessageAddressTableLookup + TransactionConfig V1 (priority_fee/compute_unit_limit/loaded_accounts_data_size_limit/heap_size) + TransactionStatusMeta + TransactionError + InnerInstructions / InnerInstruction + TokenBalance / UiTokenAmount + ReturnData / Reward + loaded writable/readonly addresses + compute_units_consumed? / cost_units? +``` + +Les marqueurs `inner_instructions_none`, `log_messages_none` et `return_data_none` sont conservés séparément de leurs payloads. Les erreurs transaction sont des bytes opaques bornés : aucun décodage Program/runtime arbitraire n'est introduit. Les types protobuf upstream restent internes. Les conversions request et les décodeurs update restent `#[cfg(test)]` jusqu'à l'ouverture du stream runtime en `pre.009`. + +OUT de `pre.007` : Blocks/block_meta/entry, stream bidi/Ping-Pong, reconnect/replay, Config V3, PublicNode et `SubscribeDeshred`. + +**Gate candidat :** audit statique clean ; compilation/Clippy/tests opérateur requis avant fermeture. diff --git a/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md b/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md index 57eac00..f848284 100644 --- a/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md +++ b/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md @@ -3,7 +3,7 @@ # Validation `0.2.9` — moteur Yellowstone + standard Solana + PublicNode -> **Statut : `0.2.9-pre.005-fix.001` est fermé sur gate opérateur intégralement vert et sans warning : fmt/audit/check/Clippy/workspace PASS, Transport 364 unit + 45 public API + 38 completeness + 4 doctests. `pre.006` est une tranche structurelle de namespace privé : `client/executor/pool/resilience/settings` deviennent `http_*` avec leurs tests miroirs, sans changement d’API publique. Les `rpc_*`, `json_rpc`, `constants` et `error` restent volontairement non préfixés lorsqu’ils sont partagés ou déjà groupés. Le forecast fonctionnel est décalé d’un numéro : Transactions `pre.007`, Blocks `pre.008`, bidi `pre.009`, reconnect `pre.010`, Config V3 `pre.011`, PublicNode final `pre.012`.** +> **Statut : `pre.006` est fermée sur gate opérateur intégralement vert : fmt/audit/check/Clippy/workspace PASS, Transport 364 unit + 45 public API + 39 completeness + 4 doctests. `0.2.9-pre.007` est candidate Transactions + `transaction_status` : filtre current complet, signature Base58 -> 64 octets, Cuckoo/token expansion, updates transaction/status et projection typée complète du `solana-storage.proto` courant incluant Transaction V1 `TransactionConfig`. Blocks restent `pre.008`; bidi reste `pre.009`.** ## 1. Autorités du gate @@ -167,6 +167,7 @@ Transaction Message MessageHeader MessageAddressTableLookup +TransactionConfig TransactionStatusMeta TransactionError InnerInstructions / InnerInstruction @@ -407,8 +408,8 @@ pre.002 DONE moteur: deps/settings/errors/channel 15 pre.003 DONE TLS/metadata + fixture + 7 unary standard 15–20 min ; gate final fix.001 PASS pre.004 DONE standard: Subscribe common/from_slot/bounds 15–20 min ; gate final fix.001 PASS pre.005 DONE standard: accounts + slots 15–20 min ; gate final fix.001 PASS -pre.006 CANDIDATE structure: namespace privé HTTP `http_*` 15–20 min -pre.007 TODO standard: transactions + transaction_status 15–20 min +pre.006 DONE structure: namespace privé HTTP `http_*` 15–20 min ; gate PASS +pre.007 CANDIDATE standard: transactions + transaction_status 15–20 min pre.008 TODO standard: blocks + block_meta + entry 15–20 min pre.009 TODO moteur: bidi/backpressure/half-close/shutdown 15–20 min pre.010 TODO moteur: reconnect/replay/gap/duplicate 15–20 min @@ -760,3 +761,55 @@ Le renommage est limité aux cinq modules dont l'ownership HTTP est sans ambigu **Verdict `pre.006` : candidate structurelle prête ; fermeture après suppressions opérateur et gate Cargo complet.** + +## 22. Gate final `pre.006` + +| Gate | Résultat final | +|------------------------------------------|----------------| +| `cargo fmt --all` | PASS | +| audit Rust workspace | PASS / clean | +| `cargo check --workspace` | PASS | +| `cargo clippy --workspace --all-targets` | PASS | +| Transport unit | 364/364 PASS | +| Transport `public_api` | 45/45 PASS | +| Transport `release_completeness` | 39/39 PASS | +| Transport doctests | 4/4 PASS | +| Core dependency canary | 3/3 PASS | +| `cargo test --workspace` | PASS | + +**Verdict : `pre.006` fermée.** + +## 23. Gate `pre.007` — Transactions + `transaction_status` candidate + +| Surface | État candidate | +|--------------------------------------------------------------------|----------------| +| workspace version | `0.2.9-pre.7` | +| `vote?` / `failed?` | SOURCE+TEST | +| signature Base58 décodée exactement 64 octets | SOURCE+TEST | +| include/exclude/required ordonnés et bornés | SOURCE+TEST | +| Cuckoo transaction include | SOURCE+TEST | +| token expansion ALL/BALANCE_CHANGED | SOURCE+TEST | +| `SubscribeUpdateTransaction` | SOURCE+TEST | +| `SubscribeUpdateTransactionStatus` | SOURCE+TEST | +| Transaction/Message/Header | SOURCE+TEST | +| compiled instructions / ALT lookups | SOURCE+TEST | +| Transaction V1 `TransactionConfig` | SOURCE+TEST | +| TransactionStatusMeta complet | SOURCE+TEST | +| inner instructions + legacy none marker | SOURCE+TEST | +| logs + legacy none marker | SOURCE+TEST | +| token balances / UiTokenAmount | SOURCE+TEST | +| loaded addresses | SOURCE+TEST | +| ReturnData + legacy none marker | SOURCE+TEST | +| rewards | SOURCE+TEST | +| compute_units_consumed / cost_units | SOURCE+TEST | +| opaque TransactionError borné | SOURCE+TEST | +| Debug sans signatures/token-balance text/payloads/logs/error bytes | SOURCE+TEST | +| Blocks / block_meta / entry | OUT pre.007 | +| stream bidi / lifecycle | OUT pre.007 | +| PublicNode / Config V3 | OUT pre.007 | +| audit Rust workspace local | PASS / clean | +| fmt/check/Clippy/tests | opérateur TODO | + +Le proto 12.6.0 ajoute à la représentation de transaction le `Message.config` optionnel pour Transaction V1 / SIMD-0385 ; la candidate le conserve explicitement. Les marqueurs legacy `inner_instructions_none`, `log_messages_none` et `return_data_none` ne sont pas fusionnés avec leurs collections/messages. + +**Verdict `pre.007` : candidate source prête ; fermeture après gate Cargo opérateur.**