v0.3.14-pre.010

This commit is contained in:
2026-09-12 14:20:02 +02:00
parent 792560191c
commit 8c4b458724
6 changed files with 620 additions and 11 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs
// version: 38
// version: 39
use sha2::Digest; // rust-rules: trait-import
@@ -18,6 +18,7 @@ pub const MIN_RAW_TRANSACTION_INGEST_HTTP_POLL_INTERVAL: std::time::Duration = s
/// Minimum number of confirmed blocks accepted during one HTTP live polling cycle.
pub const MIN_RAW_TRANSACTION_INGEST_HTTP_POLL_MAX_BLOCKS_PER_CYCLE: u16 = 1;
const MAX_RAW_TRANSACTION_INGEST_REPAIR_BURST: usize = 1;
const RAW_TRANSACTION_INGEST_HELIUS_TRANSACTION_FILTER_FINGERPRINT_DOMAIN: &[u8] = b"ksp.raw_transaction_ingest.helius_transaction.filter.v1\0";
const RAW_TRANSACTION_INGEST_HELIUS_TRANSACTION_HTTP_PROTOCOL: &str = "helius_ws_http";
const RAW_TRANSACTION_INGEST_HELIUS_TRANSACTION_HTTP_SOURCE_KEY_DOMAIN: &[u8] = b"ksp.raw_transaction_ingest.helius_transaction_http.source_key.v1\0";
@@ -36,6 +37,197 @@ const RAW_TRANSACTION_INGEST_YELLOWSTONE_FILTER_FINGERPRINT_DOMAIN: &[u8] = b"ks
const RAW_TRANSACTION_INGEST_YELLOWSTONE_HTTP_PROTOCOL: &str = "yellowstone_http";
const RAW_TRANSACTION_INGEST_YELLOWSTONE_HTTP_SOURCE_KEY_DOMAIN: &[u8] = b"ksp.raw_transaction_ingest.yellowstone_http.source_key.v1\0";
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum RawTransactionIngestTrafficClass {
Nominal,
Repair,
}
struct RawTransactionIngestFairTurnState {
active: bool,
last_granted: std::option::Option<RawTransactionIngestTrafficClass>,
next_ticket: u64,
waiters: std::collections::VecDeque<(u64, RawTransactionIngestTrafficClass)>,
}
struct RawTransactionIngestFairTurnGate {
max_waiters: usize,
notify: tokio::sync::Notify,
state: std::sync::Mutex<RawTransactionIngestFairTurnState>,
}
struct RawTransactionIngestFairWaiterGuard {
armed: bool,
gate: std::sync::Arc<RawTransactionIngestFairTurnGate>,
ticket: u64,
}
struct RawTransactionIngestFairTurnGuard {
gate: std::sync::Arc<RawTransactionIngestFairTurnGate>,
}
impl RawTransactionIngestFairTurnGate {
fn new(max_waiters: usize) -> Self {
return Self {
max_waiters,
notify: tokio::sync::Notify::new(),
state: std::sync::Mutex::new(RawTransactionIngestFairTurnState {
active: false,
last_granted: std::option::Option::None,
next_ticket: 1,
waiters: std::collections::VecDeque::new(),
}),
};
}
async fn acquire(self: std::sync::Arc<Self>, class: RawTransactionIngestTrafficClass) -> ksp_core_lib::Result<RawTransactionIngestFairTurnGuard> {
let mut waiter = match RawTransactionIngestFairWaiterGuard::register(std::sync::Arc::clone(&self), class) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
loop {
let notify_gate = std::sync::Arc::clone(&self);
let notified = notify_gate.notify.notified();
let granted = {
let mut state = match self.state.lock() {
std::result::Result::Ok(value) => value,
std::result::Result::Err(poisoned) => poisoned.into_inner(),
};
if state.active || fair_turn_selected_ticket(&state) != std::option::Option::Some(waiter.ticket) {
false
} else {
let position = state.waiters.iter().position(|(ticket, _)| return *ticket == waiter.ticket);
let position = match position {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("source.fairness_waiter_missing")),
};
let removed = state.waiters.remove(position);
let removed = match removed {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("source.fairness_waiter_missing")),
};
if removed.1 != class {
return std::result::Result::Err(crate::runtime_error("source.fairness_waiter_class_mismatch"));
}
state.active = true;
state.last_granted = std::option::Option::Some(class);
true
}
};
if granted {
waiter.armed = false;
std::mem::drop(notified);
std::mem::drop(notify_gate);
return std::result::Result::Ok(RawTransactionIngestFairTurnGuard { gate: self });
}
notified.await;
}
}
}
impl RawTransactionIngestFairWaiterGuard {
fn register(gate: std::sync::Arc<RawTransactionIngestFairTurnGate>, class: RawTransactionIngestTrafficClass) -> ksp_core_lib::Result<Self> {
let ticket = {
let mut state = match gate.state.lock() {
std::result::Result::Ok(value) => value,
std::result::Result::Err(poisoned) => poisoned.into_inner(),
};
if gate.max_waiters == 0 || state.waiters.len() >= gate.max_waiters {
return std::result::Result::Err(crate::runtime_error("source.fairness_waiter_capacity_exceeded"));
}
let ticket = state.next_ticket;
state.next_ticket = match state.next_ticket.checked_add(1) {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("source.fairness_ticket_exhausted")),
};
state.waiters.push_back((ticket, class));
ticket
};
gate.notify.notify_waiters();
return std::result::Result::Ok(Self { armed: true, gate, ticket });
}
}
impl std::ops::Drop for RawTransactionIngestFairWaiterGuard {
fn drop(&mut self) {
if !self.armed {
return;
}
let mut state = match self.gate.state.lock() {
std::result::Result::Ok(value) => value,
std::result::Result::Err(poisoned) => poisoned.into_inner(),
};
if let std::option::Option::Some(position) = state.waiters.iter().position(|(ticket, _)| return *ticket == self.ticket) {
let _removed = state.waiters.remove(position);
}
std::mem::drop(state);
self.gate.notify.notify_waiters();
return;
}
}
impl std::ops::Drop for RawTransactionIngestFairTurnGuard {
fn drop(&mut self) {
let mut state = match self.gate.state.lock() {
std::result::Result::Ok(value) => value,
std::result::Result::Err(poisoned) => poisoned.into_inner(),
};
state.active = false;
std::mem::drop(state);
self.gate.notify.notify_waiters();
return;
}
}
fn fair_turn_selected_ticket(state: &RawTransactionIngestFairTurnState) -> std::option::Option<u64> {
let nominal = state.waiters.iter().find(|(_, class)| return *class == RawTransactionIngestTrafficClass::Nominal);
let repair = state.waiters.iter().find(|(_, class)| return *class == RawTransactionIngestTrafficClass::Repair);
let selected = match (nominal, repair) {
(std::option::Option::Some(nominal), std::option::Option::Some(repair)) => match state.last_granted {
std::option::Option::Some(RawTransactionIngestTrafficClass::Nominal) => repair,
std::option::Option::Some(RawTransactionIngestTrafficClass::Repair) | std::option::Option::None => nominal,
},
(std::option::Option::Some(nominal), std::option::Option::None) => nominal,
(std::option::Option::None, std::option::Option::Some(repair)) => repair,
(std::option::Option::None, std::option::Option::None) => return std::option::Option::None,
};
return std::option::Option::Some(selected.0);
}
fn validate_repair_fairness_contract(admission_capacity: usize, persistence_concurrency: usize) -> ksp_core_lib::Result<()> {
if admission_capacity == 0
|| persistence_concurrency == 0
|| repair_block_fetch_limit(persistence_concurrency) == 0
|| MAX_RAW_TRANSACTION_INGEST_REPAIR_BURST != 1
|| !repair_fairness_catalog_is_complete()
{
return std::result::Result::Err(crate::runtime_error("runtime_resources.repair_fairness_invalid"));
}
return std::result::Result::Ok(());
}
fn repair_block_fetch_limit(existing_capacity: usize) -> usize {
return existing_capacity.min(4);
}
fn repair_fairness_catalog_is_complete() -> bool {
let mut state = RawTransactionIngestFairTurnState {
active: false,
last_granted: std::option::Option::None,
next_ticket: 3,
waiters: std::collections::VecDeque::from([(1, RawTransactionIngestTrafficClass::Nominal), (2, RawTransactionIngestTrafficClass::Repair)]),
};
if fair_turn_selected_ticket(&state) != std::option::Option::Some(1) {
return false;
}
state.last_granted = std::option::Option::Some(RawTransactionIngestTrafficClass::Nominal);
if fair_turn_selected_ticket(&state) != std::option::Option::Some(2) {
return false;
}
state.last_granted = std::option::Option::Some(RawTransactionIngestTrafficClass::Repair);
return fair_turn_selected_ticket(&state) == std::option::Option::Some(1);
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum RawTransactionIngestHttpDiscoveryStrategy {
ClosedRange,
@@ -2336,6 +2528,9 @@ impl crate::RawTransactionIngestRuntimeResources {
{
return std::result::Result::Err(error);
}
if let std::result::Result::Err(error) = validate_repair_fairness_contract(settings.admission_queue_capacity(), settings.persistence_concurrency()) {
return std::result::Result::Err(error);
}
let hydration_pending_budget = settings.admission_queue_capacity();
let hydration_in_flight_budget = settings.persistence_concurrency();
let inventory = std::sync::Arc::new(std::sync::Mutex::new(RawTransactionIngestSourceInventory::new(source_keys)));
@@ -4263,6 +4458,7 @@ enum RawTransactionIngestSharedHydrationResult {
}
struct RawTransactionIngestGlobalHydrationRegistry {
hydration_fairness: std::sync::Arc<RawTransactionIngestFairTurnGate>,
hydration_permits: std::sync::Arc<tokio::sync::Semaphore>,
max_pending: usize,
pending: std::sync::Mutex<
@@ -4276,12 +4472,26 @@ struct RawTransactionIngestGlobalHydrationRegistry {
impl RawTransactionIngestGlobalHydrationRegistry {
fn new(max_pending: usize, max_in_flight: usize) -> Self {
return Self {
hydration_fairness: std::sync::Arc::new(RawTransactionIngestFairTurnGate::new(max_pending)),
hydration_permits: std::sync::Arc::new(tokio::sync::Semaphore::new(max_in_flight)),
max_pending,
pending: std::sync::Mutex::new(std::collections::BTreeMap::new()),
};
}
async fn acquire_hydration_permit(&self, class: RawTransactionIngestTrafficClass) -> ksp_core_lib::Result<tokio::sync::OwnedSemaphorePermit> {
let turn = match std::sync::Arc::clone(&self.hydration_fairness).acquire(class).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let permit = std::sync::Arc::clone(&self.hydration_permits).acquire_owned().await;
std::mem::drop(turn);
return match permit {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(crate::runtime_error("source.global_hydration_permit_closed")),
};
}
fn subscribe_or_lead(
&self,
key: &RawTransactionIngestHydrationKey,
@@ -4599,10 +4809,9 @@ async fn fetch_hydration_shared(
};
if leader {
let mut leader_guard = RawTransactionIngestHydrationLeaderGuard::new(std::sync::Arc::clone(&global_registry), key.clone());
let permit = std::sync::Arc::clone(&global_registry.hydration_permits).acquire_owned().await;
let _permit = match permit {
let _permit = match global_registry.acquire_hydration_permit(RawTransactionIngestTrafficClass::Nominal).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(crate::runtime_error("source.global_hydration_permit_closed")),
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let fetched = fetch_hydration(http_pool, hydration_role, expected_network, key.clone(), commitment).await;
return match fetched {