v0.3.13-pre.008
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs
|
||||
// version: 23
|
||||
// version: 24
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -124,12 +124,16 @@ pub(crate) use self::error::store_error;
|
||||
pub(crate) use self::persistence::RawTransactionIngestEntityPersistence;
|
||||
/// Observation disposition produced by one successful Worker Store persistence attempt.
|
||||
pub(crate) use self::persistence::RawTransactionIngestObservationPersistence;
|
||||
/// Private bounded run-local cache serializing repeated canonical identities before Store writes.
|
||||
pub(crate) use self::persistence::RawTransactionIngestPersistenceConvergence;
|
||||
/// Classified successful outcome of one atomic Worker Store persistence attempt.
|
||||
pub(crate) use self::persistence::RawTransactionIngestPersistenceOutcome;
|
||||
/// Private backend-neutral Store persistence port used by the Worker and deterministic tests.
|
||||
pub(crate) use self::persistence::RawTransactionIngestPersistencePort;
|
||||
/// Persists one already-canonical Worker acquisition through the private Store port in `Normal` mode.
|
||||
pub(crate) use self::persistence::persist_raw_transaction_ingest_acquisition;
|
||||
/// Persists one canonical acquisition through the bounded cross-source convergence cache.
|
||||
pub(crate) use self::persistence::persist_raw_transaction_ingest_converged_acquisition;
|
||||
/// Private latest-value processing-frontier projection emitted by the productive source task.
|
||||
pub(crate) use self::snapshot::RawTransactionIngestProcessingFrontierProjection;
|
||||
/// Private latest-value publisher and checked counter owner shared by the Worker supervisor.
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs
|
||||
// version: 1
|
||||
// version: 2
|
||||
|
||||
/// Canonical entity disposition produced by one Worker Store persistence attempt.
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
@@ -56,6 +56,12 @@ pub(crate) trait RawTransactionIngestPersistencePort: std::marker::Send + std::m
|
||||
observation: ksp_store_lib::RawTransactionObservation,
|
||||
mode: ksp_store_lib::RawTransactionAcquisitionMode,
|
||||
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawAcquisitionWriteOutcome>>;
|
||||
|
||||
/// Records one additional observation for an already durable canonical RAW transaction.
|
||||
fn record_observation<'a>(
|
||||
&'a self,
|
||||
observation: ksp_store_lib::RawTransactionObservation,
|
||||
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawObservationWriteOutcome>>;
|
||||
}
|
||||
|
||||
impl crate::RawTransactionIngestPersistencePort for ksp_store_lib::Store {
|
||||
@@ -71,6 +77,65 @@ impl crate::RawTransactionIngestPersistencePort for ksp_store_lib::Store {
|
||||
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawAcquisitionWriteOutcome>> {
|
||||
return ksp_store_lib::RawTransactionWrite::persist_raw_transaction_acquisition(self, transaction, observation, mode);
|
||||
}
|
||||
|
||||
fn record_observation<'a>(
|
||||
&'a self,
|
||||
observation: ksp_store_lib::RawTransactionObservation,
|
||||
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawObservationWriteOutcome>> {
|
||||
return ksp_store_lib::RawTransactionObservationWrite::record_raw_transaction_observation(self, observation);
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Eq, Ord, PartialEq, PartialOrd)]
|
||||
struct RawTransactionIngestPersistenceKey {
|
||||
network: ksp_store_lib::RawNetworkId,
|
||||
signature: ksp_store_lib::RawTransactionSignature,
|
||||
}
|
||||
|
||||
/// Private bounded run-local cache serializing repeated canonical identities before Store writes.
|
||||
pub(crate) struct RawTransactionIngestPersistenceConvergence {
|
||||
max_entries: usize,
|
||||
entries: std::sync::Mutex<
|
||||
std::collections::BTreeMap<RawTransactionIngestPersistenceKey, std::sync::Arc<tokio::sync::Mutex<std::option::Option<ksp_store_lib::RawContentHash>>>>,
|
||||
>,
|
||||
}
|
||||
|
||||
impl crate::RawTransactionIngestPersistenceConvergence {
|
||||
/// Creates one bounded run-local convergence cache using the caller-provided effective runtime bound.
|
||||
#[must_use]
|
||||
pub(crate) fn new(max_entries: usize) -> Self {
|
||||
return Self { max_entries, entries: std::sync::Mutex::new(std::collections::BTreeMap::new()) };
|
||||
}
|
||||
|
||||
fn entry(
|
||||
&self,
|
||||
key: RawTransactionIngestPersistenceKey,
|
||||
) -> std::option::Option<std::sync::Arc<tokio::sync::Mutex<std::option::Option<ksp_store_lib::RawContentHash>>>> {
|
||||
let mut entries = match self.entries.lock() {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(poisoned) => poisoned.into_inner(),
|
||||
};
|
||||
if let std::option::Option::Some(entry) = entries.get(&key) {
|
||||
return std::option::Option::Some(std::sync::Arc::clone(entry));
|
||||
}
|
||||
if entries.len() >= self.max_entries {
|
||||
let removable = entries.iter().find_map(|(candidate, entry)| {
|
||||
if std::sync::Arc::strong_count(entry) == 1 {
|
||||
return std::option::Option::Some(candidate.clone());
|
||||
}
|
||||
return std::option::Option::None;
|
||||
});
|
||||
if let std::option::Option::Some(removable) = removable {
|
||||
entries.remove(&removable);
|
||||
}
|
||||
}
|
||||
if entries.len() >= self.max_entries {
|
||||
return std::option::Option::None;
|
||||
}
|
||||
let entry = std::sync::Arc::new(tokio::sync::Mutex::new(std::option::Option::None));
|
||||
entries.insert(key, std::sync::Arc::clone(&entry));
|
||||
return std::option::Option::Some(entry);
|
||||
}
|
||||
}
|
||||
|
||||
/// Persists one already-canonical Worker acquisition through the private Store port in `Normal` mode.
|
||||
@@ -96,6 +161,52 @@ where
|
||||
};
|
||||
}
|
||||
|
||||
/// Persists one canonical acquisition through the bounded cross-source convergence cache.
|
||||
pub(crate) async fn persist_raw_transaction_ingest_converged_acquisition<P>(
|
||||
port: &P,
|
||||
acquisition: ksp_raw_transaction_lib::RawTransactionAcquisition,
|
||||
convergence: &crate::RawTransactionIngestPersistenceConvergence,
|
||||
) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome>
|
||||
where
|
||||
P: crate::RawTransactionIngestPersistencePort + ?Sized,
|
||||
{
|
||||
let reference = acquisition.transaction().reference().clone();
|
||||
if !port.network_matches(reference.network()) {
|
||||
return std::result::Result::Err(crate::runtime_error("persistence.store_network_mismatch"));
|
||||
}
|
||||
if acquisition.observation().transaction() != &reference {
|
||||
return std::result::Result::Err(crate::runtime_error("persistence.acquisition_reference_mismatch"));
|
||||
}
|
||||
let content_hash = acquisition.transaction().payload().content_hash();
|
||||
let key = RawTransactionIngestPersistenceKey { network: reference.network().clone(), signature: reference.signature() };
|
||||
let entry = convergence.entry(key);
|
||||
let entry = match entry {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return persist_raw_transaction_ingest_acquisition(port, acquisition).await,
|
||||
};
|
||||
let mut known_hash = entry.lock().await;
|
||||
if let std::option::Option::Some(known_hash_value) = *known_hash {
|
||||
if known_hash_value != content_hash {
|
||||
return std::result::Result::Err(crate::content_conflict_error());
|
||||
}
|
||||
let (_transaction, observation) = acquisition.into_parts();
|
||||
let result = port.record_observation(observation).await;
|
||||
return match result {
|
||||
std::result::Result::Ok(outcome) => map_additional_observation_outcome(outcome),
|
||||
std::result::Result::Err(error) => map_store_error(error),
|
||||
};
|
||||
}
|
||||
let outcome = persist_raw_transaction_ingest_acquisition(port, acquisition).await;
|
||||
let outcome = match outcome {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
if outcome.entity() != crate::RawTransactionIngestEntityPersistence::SkippedPurged {
|
||||
*known_hash = std::option::Option::Some(content_hash);
|
||||
}
|
||||
return std::result::Result::Ok(outcome);
|
||||
}
|
||||
|
||||
fn map_store_error(error: ksp_core_lib::Error) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome> {
|
||||
if error.code() == ksp_store_lib::ERROR_CODE_RAW_CONFLICT {
|
||||
return std::result::Result::Err(crate::content_conflict_error());
|
||||
@@ -133,6 +244,24 @@ fn map_store_outcome(outcome: ksp_store_lib::RawAcquisitionWriteOutcome) -> ksp_
|
||||
return std::result::Result::Err(crate::runtime_error("persistence.store_outcome_invalid"));
|
||||
}
|
||||
|
||||
fn map_additional_observation_outcome(
|
||||
outcome: ksp_store_lib::RawObservationWriteOutcome,
|
||||
) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome> {
|
||||
return match outcome {
|
||||
ksp_store_lib::RawObservationWriteOutcome::Inserted => std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
|
||||
entity: crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
|
||||
observation: crate::RawTransactionIngestObservationPersistence::Inserted,
|
||||
}),
|
||||
ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent => std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
|
||||
entity: crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
|
||||
observation: crate::RawTransactionIngestObservationPersistence::AlreadyPresent,
|
||||
}),
|
||||
ksp_store_lib::RawObservationWriteOutcome::NotRecorded => {
|
||||
std::result::Result::Err(crate::runtime_error("persistence.additional_observation_not_recorded"))
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "../unit_tests/persistence.rs"]
|
||||
mod tests;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs
|
||||
// version: 14
|
||||
// version: 15
|
||||
|
||||
type PersistencePort = std::sync::Arc<dyn crate::RawTransactionIngestPersistencePort + 'static>;
|
||||
type PersistenceTasks = tokio::task::JoinSet<ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome>>;
|
||||
@@ -189,6 +189,7 @@ async fn drain_admission_and_persistence(
|
||||
lifecycle: &ksp_worker_api::WorkerLifecycle,
|
||||
admission: &mut crate::RawTransactionAdmission,
|
||||
persistence: &mut PersistenceTasks,
|
||||
persistence_convergence: &std::sync::Arc<crate::RawTransactionIngestPersistenceConvergence>,
|
||||
port: &std::option::Option<PersistencePort>,
|
||||
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
|
||||
) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
@@ -207,7 +208,7 @@ async fn drain_admission_and_persistence(
|
||||
let backpressure_wait_observed = admission.take_backpressure_wait_observed();
|
||||
match received {
|
||||
std::result::Result::Ok(std::option::Option::Some(acquisition)) => {
|
||||
let spawned = spawn_persistence(persistence, port, acquisition);
|
||||
let spawned = spawn_persistence(persistence, persistence_convergence, port, acquisition);
|
||||
let published = snapshots.record_admission_success(lifecycle.state(), admission.queue_depth(), persistence.len(), backpressure_wait_observed);
|
||||
if !spawned && fault.is_none() {
|
||||
fault = std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
||||
@@ -247,11 +248,12 @@ async fn drain_owned_work(
|
||||
children: &mut SourceTasks,
|
||||
admission: &mut crate::RawTransactionAdmission,
|
||||
persistence: &mut PersistenceTasks,
|
||||
persistence_convergence: &std::sync::Arc<crate::RawTransactionIngestPersistenceConvergence>,
|
||||
port: &std::option::Option<PersistencePort>,
|
||||
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
|
||||
) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
let drain = async {
|
||||
let mut fault = drain_admission_and_persistence(settings, lifecycle, admission, persistence, port, snapshots).await;
|
||||
let mut fault = drain_admission_and_persistence(settings, lifecycle, admission, persistence, persistence_convergence, port, snapshots).await;
|
||||
let child_fault = drain_children(lifecycle.state(), children, snapshots).await;
|
||||
fault = merge_fault(fault, child_fault);
|
||||
return fault;
|
||||
@@ -380,6 +382,9 @@ async fn run_supervisor<Spawner>(
|
||||
}
|
||||
let mut children = SourceTasks::new();
|
||||
let mut persistence = PersistenceTasks::new();
|
||||
let persistence_convergence = std::sync::Arc::new(crate::RawTransactionIngestPersistenceConvergence::new(
|
||||
settings.admission_queue_capacity().max(settings.persistence_concurrency()),
|
||||
));
|
||||
let (source_stop_sender, source_stop_receiver) = tokio::sync::watch::channel(false);
|
||||
let (mut admission, admission_sender) = crate::RawTransactionAdmission::new(settings.admission_queue_capacity());
|
||||
source_spawner(&mut children, source_stop_receiver, admission_sender);
|
||||
@@ -391,6 +396,7 @@ async fn run_supervisor<Spawner>(
|
||||
&mut children,
|
||||
&mut admission,
|
||||
&mut persistence,
|
||||
&persistence_convergence,
|
||||
&port,
|
||||
&mut snapshots,
|
||||
&mut processing_frontier_receiver,
|
||||
@@ -399,7 +405,8 @@ async fn run_supervisor<Spawner>(
|
||||
source_stop_sender.send_replace(true);
|
||||
let stopping_fault = begin_stopping(&mut lifecycle, &mut snapshots, admission.queue_depth(), persistence.len());
|
||||
fault = merge_fault(fault, stopping_fault);
|
||||
let drain_fault = drain_owned_work(&settings, &lifecycle, &mut children, &mut admission, &mut persistence, &port, &mut snapshots).await;
|
||||
let drain_fault =
|
||||
drain_owned_work(&settings, &lifecycle, &mut children, &mut admission, &mut persistence, &persistence_convergence, &port, &mut snapshots).await;
|
||||
match drain_fault {
|
||||
std::option::Option::Some(code) if code == crate::ERROR_CODE_RAW_TRANSACTION_INGEST_DRAIN_TIMEOUT => {
|
||||
fault = std::option::Option::Some(code);
|
||||
@@ -418,6 +425,7 @@ async fn run_supervisor<Spawner>(
|
||||
|
||||
fn spawn_persistence(
|
||||
persistence: &mut PersistenceTasks,
|
||||
persistence_convergence: &std::sync::Arc<crate::RawTransactionIngestPersistenceConvergence>,
|
||||
port: &std::option::Option<PersistencePort>,
|
||||
acquisition: ksp_raw_transaction_lib::RawTransactionAcquisition,
|
||||
) -> bool {
|
||||
@@ -425,8 +433,9 @@ fn spawn_persistence(
|
||||
std::option::Option::Some(value) => std::sync::Arc::clone(value),
|
||||
std::option::Option::None => return false,
|
||||
};
|
||||
let persistence_convergence = std::sync::Arc::clone(persistence_convergence);
|
||||
let _abort_handle = persistence.spawn(async move {
|
||||
return crate::persist_raw_transaction_ingest_acquisition(port.as_ref(), acquisition).await;
|
||||
return crate::persist_raw_transaction_ingest_converged_acquisition(port.as_ref(), acquisition, persistence_convergence.as_ref()).await;
|
||||
});
|
||||
return true;
|
||||
}
|
||||
@@ -507,6 +516,7 @@ async fn supervise_until_stop(
|
||||
children: &mut SourceTasks,
|
||||
admission: &mut crate::RawTransactionAdmission,
|
||||
persistence: &mut PersistenceTasks,
|
||||
persistence_convergence: &std::sync::Arc<crate::RawTransactionIngestPersistenceConvergence>,
|
||||
port: &std::option::Option<PersistencePort>,
|
||||
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
|
||||
processing_frontier_receiver: &mut std::option::Option<ProcessingFrontierReceiver>,
|
||||
@@ -578,7 +588,7 @@ async fn supervise_until_stop(
|
||||
match received {
|
||||
std::result::Result::Ok(std::option::Option::Some(acquisition)) => {
|
||||
let backpressure_wait_observed = admission.take_backpressure_wait_observed();
|
||||
let spawned = spawn_persistence(persistence, port, acquisition);
|
||||
let spawned = spawn_persistence(persistence, persistence_convergence, port, acquisition);
|
||||
let published = snapshots.record_admission_success(
|
||||
lifecycle.state(),
|
||||
admission.queue_depth(),
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs
|
||||
// version: 21
|
||||
// version: 22
|
||||
|
||||
use sha2::Digest; // rust-rules: trait-import
|
||||
|
||||
@@ -90,17 +90,20 @@ impl RawTransactionIngestLiveSource {
|
||||
settings: crate::RawTransactionIngestSettings,
|
||||
stop_receiver: tokio::sync::watch::Receiver<bool>,
|
||||
admission_sender: tokio::sync::mpsc::Sender<crate::RawTransactionIngress>,
|
||||
global_hydration_registry: std::sync::Arc<RawTransactionIngestGlobalHydrationRegistry>,
|
||||
inventory_publisher: RawTransactionIngestSourceInventoryPublisher,
|
||||
) -> ksp_core_lib::Result<()> {
|
||||
let (source_frontier_sender, mut source_frontier_receiver) =
|
||||
tokio::sync::watch::channel(crate::RawTransactionIngestProcessingFrontierProjection::empty());
|
||||
let mut source_future = std::boxed::Box::pin(async move {
|
||||
return match self {
|
||||
Self::HeliusTransaction(source) => source.run(settings, stop_receiver, admission_sender, source_frontier_sender).await,
|
||||
Self::HeliusTransaction(source) => {
|
||||
source.run(settings, stop_receiver, admission_sender, source_frontier_sender, global_hydration_registry).await
|
||||
},
|
||||
Self::HttpBlockPolling(source) => source.run(settings, stop_receiver, admission_sender, source_frontier_sender).await,
|
||||
Self::StandardBlock(source) => source.run(settings, stop_receiver, admission_sender, source_frontier_sender).await,
|
||||
Self::StandardLogs(source) => source.run(settings, stop_receiver, admission_sender, source_frontier_sender).await,
|
||||
Self::Yellowstone(source) => source.run(settings, stop_receiver, admission_sender, source_frontier_sender).await,
|
||||
Self::StandardLogs(source) => source.run(settings, stop_receiver, admission_sender, source_frontier_sender, global_hydration_registry).await,
|
||||
Self::Yellowstone(source) => source.run(settings, stop_receiver, admission_sender, source_frontier_sender, global_hydration_registry).await,
|
||||
};
|
||||
});
|
||||
loop {
|
||||
@@ -624,6 +627,7 @@ impl crate::RawTransactionIngestYellowstoneSource {
|
||||
mut stop_receiver: tokio::sync::watch::Receiver<bool>,
|
||||
admission_sender: tokio::sync::mpsc::Sender<crate::RawTransactionIngress>,
|
||||
processing_frontier_sender: tokio::sync::watch::Sender<crate::RawTransactionIngestProcessingFrontierProjection>,
|
||||
global_hydration_registry: std::sync::Arc<RawTransactionIngestGlobalHydrationRegistry>,
|
||||
) -> ksp_core_lib::Result<()> {
|
||||
let opened = tokio::select! {
|
||||
biased;
|
||||
@@ -637,7 +641,7 @@ impl crate::RawTransactionIngestYellowstoneSource {
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(source_transport_error(error.code())),
|
||||
};
|
||||
let hydration = self.hydration_context();
|
||||
let mut coordinator = RawTransactionIngestHydrationCoordinator::new(&settings);
|
||||
let mut coordinator = RawTransactionIngestHydrationCoordinator::with_global_registry(&settings, global_hydration_registry);
|
||||
let mut processing_frontier = RawTransactionIngestProcessingFrontierReporter::new(processing_frontier_sender);
|
||||
let snapshot_source = session.snapshot_source();
|
||||
if let std::result::Result::Err(error) = processing_frontier.observe_session_snapshot(snapshot_source.current()) {
|
||||
@@ -860,6 +864,7 @@ impl crate::RawTransactionIngestHeliusTransactionSource {
|
||||
mut stop_receiver: tokio::sync::watch::Receiver<bool>,
|
||||
admission_sender: tokio::sync::mpsc::Sender<crate::RawTransactionIngress>,
|
||||
processing_frontier_sender: tokio::sync::watch::Sender<crate::RawTransactionIngestProcessingFrontierProjection>,
|
||||
global_hydration_registry: std::sync::Arc<RawTransactionIngestGlobalHydrationRegistry>,
|
||||
) -> ksp_core_lib::Result<()> {
|
||||
let connected = tokio::select! {
|
||||
biased;
|
||||
@@ -889,7 +894,7 @@ impl crate::RawTransactionIngestHeliusTransactionSource {
|
||||
},
|
||||
};
|
||||
let hydration = self.hydration_context();
|
||||
let mut coordinator = RawTransactionIngestHydrationCoordinator::new(&settings);
|
||||
let mut coordinator = RawTransactionIngestHydrationCoordinator::with_global_registry(&settings, global_hydration_registry);
|
||||
let mut processing_frontier = RawTransactionIngestProcessingFrontierReporter::new(processing_frontier_sender);
|
||||
processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Active);
|
||||
let mut fault = std::option::Option::None;
|
||||
@@ -1578,6 +1583,7 @@ impl crate::RawTransactionIngestStandardLogsSource {
|
||||
mut stop_receiver: tokio::sync::watch::Receiver<bool>,
|
||||
admission_sender: tokio::sync::mpsc::Sender<crate::RawTransactionIngress>,
|
||||
processing_frontier_sender: tokio::sync::watch::Sender<crate::RawTransactionIngestProcessingFrontierProjection>,
|
||||
global_hydration_registry: std::sync::Arc<RawTransactionIngestGlobalHydrationRegistry>,
|
||||
) -> ksp_core_lib::Result<()> {
|
||||
let connected = tokio::select! {
|
||||
biased;
|
||||
@@ -1607,7 +1613,7 @@ impl crate::RawTransactionIngestStandardLogsSource {
|
||||
},
|
||||
};
|
||||
let hydration = self.hydration_context();
|
||||
let mut coordinator = RawTransactionIngestHydrationCoordinator::new(&settings);
|
||||
let mut coordinator = RawTransactionIngestHydrationCoordinator::with_global_registry(&settings, global_hydration_registry);
|
||||
let mut processing_frontier = RawTransactionIngestProcessingFrontierReporter::new(processing_frontier_sender);
|
||||
processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Active);
|
||||
let mut fault = std::option::Option::None;
|
||||
@@ -1934,6 +1940,7 @@ impl crate::RawTransactionIngestRuntimeResources {
|
||||
}
|
||||
let source_keys = self.sources.iter().map(RawTransactionIngestLiveSource::source_key).collect::<std::vec::Vec<_>>();
|
||||
let inventory = std::sync::Arc::new(std::sync::Mutex::new(RawTransactionIngestSourceInventory::new(source_keys)));
|
||||
let global_hydration_registry = std::sync::Arc::new(RawTransactionIngestGlobalHydrationRegistry::new(&settings));
|
||||
let (source_stop_sender, source_stop_receiver) = tokio::sync::watch::channel(false);
|
||||
let mut children = tokio::task::JoinSet::new();
|
||||
for (entry_index, source) in self.sources.into_iter().enumerate() {
|
||||
@@ -1946,8 +1953,9 @@ impl crate::RawTransactionIngestRuntimeResources {
|
||||
let source_admission_sender = admission_sender.clone();
|
||||
let source_settings = settings.clone();
|
||||
let source_stop_receiver = source_stop_receiver.clone();
|
||||
let source_global_hydration_registry = std::sync::Arc::clone(&global_hydration_registry);
|
||||
let _abort_handle = children.spawn(async move {
|
||||
return source.run(source_settings, source_stop_receiver, source_admission_sender, publisher).await;
|
||||
return source.run(source_settings, source_stop_receiver, source_admission_sender, source_global_hydration_registry, publisher).await;
|
||||
});
|
||||
}
|
||||
std::mem::drop(admission_sender);
|
||||
@@ -3336,9 +3344,68 @@ struct RawTransactionIngestHydrationFetch {
|
||||
observed: RawTransactionIngestObservedTransaction,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
enum RawTransactionIngestSharedHydrationResult {
|
||||
Available(RawTransactionIngestObservedTransaction),
|
||||
Failed(ksp_core_lib::ErrorCode),
|
||||
}
|
||||
|
||||
struct RawTransactionIngestGlobalHydrationRegistry {
|
||||
max_pending: usize,
|
||||
pending: std::sync::Mutex<
|
||||
std::collections::BTreeMap<
|
||||
RawTransactionIngestHydrationKey,
|
||||
tokio::sync::watch::Sender<std::option::Option<RawTransactionIngestSharedHydrationResult>>,
|
||||
>,
|
||||
>,
|
||||
}
|
||||
|
||||
impl RawTransactionIngestGlobalHydrationRegistry {
|
||||
fn new(settings: &crate::RawTransactionIngestSettings) -> Self {
|
||||
return Self {
|
||||
max_pending: settings.admission_queue_capacity(),
|
||||
pending: std::sync::Mutex::new(std::collections::BTreeMap::new()),
|
||||
};
|
||||
}
|
||||
|
||||
fn subscribe_or_lead(
|
||||
&self,
|
||||
key: &RawTransactionIngestHydrationKey,
|
||||
) -> ksp_core_lib::Result<(bool, tokio::sync::watch::Receiver<std::option::Option<RawTransactionIngestSharedHydrationResult>>)> {
|
||||
let mut pending = match self.pending.lock() {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(poisoned) => poisoned.into_inner(),
|
||||
};
|
||||
if let std::option::Option::Some(sender) = pending.get(key) {
|
||||
return std::result::Result::Ok((false, sender.subscribe()));
|
||||
}
|
||||
if pending.len() >= self.max_pending {
|
||||
return std::result::Result::Err(crate::runtime_error("source.global_hydration_pending_saturated"));
|
||||
}
|
||||
let (sender, receiver) = tokio::sync::watch::channel(std::option::Option::None);
|
||||
pending.insert(key.clone(), sender);
|
||||
return std::result::Result::Ok((true, receiver));
|
||||
}
|
||||
|
||||
fn publish_and_remove(&self, key: &RawTransactionIngestHydrationKey, result: RawTransactionIngestSharedHydrationResult) {
|
||||
let sender = {
|
||||
let mut pending = match self.pending.lock() {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(poisoned) => poisoned.into_inner(),
|
||||
};
|
||||
pending.remove(key)
|
||||
};
|
||||
if let std::option::Option::Some(sender) = sender {
|
||||
sender.send_replace(std::option::Option::Some(result));
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
type RawTransactionIngestHydrationTasks = tokio::task::JoinSet<ksp_core_lib::Result<RawTransactionIngestHydrationFetch>>;
|
||||
|
||||
struct RawTransactionIngestHydrationCoordinator {
|
||||
global_registry: std::sync::Arc<RawTransactionIngestGlobalHydrationRegistry>,
|
||||
max_in_flight: usize,
|
||||
max_pending_signals: usize,
|
||||
pending_signal_count: usize,
|
||||
@@ -3348,7 +3415,15 @@ struct RawTransactionIngestHydrationCoordinator {
|
||||
|
||||
impl RawTransactionIngestHydrationCoordinator {
|
||||
fn new(settings: &crate::RawTransactionIngestSettings) -> Self {
|
||||
return Self::with_global_registry(settings, std::sync::Arc::new(RawTransactionIngestGlobalHydrationRegistry::new(settings)));
|
||||
}
|
||||
|
||||
fn with_global_registry(
|
||||
settings: &crate::RawTransactionIngestSettings,
|
||||
global_registry: std::sync::Arc<RawTransactionIngestGlobalHydrationRegistry>,
|
||||
) -> Self {
|
||||
return Self {
|
||||
global_registry,
|
||||
max_in_flight: settings.persistence_concurrency(),
|
||||
max_pending_signals: settings.admission_queue_capacity(),
|
||||
pending_signal_count: 0,
|
||||
@@ -3419,8 +3494,9 @@ impl RawTransactionIngestHydrationCoordinator {
|
||||
let role = hydration.hydration_role.clone();
|
||||
let expected_network = hydration.network.clone();
|
||||
let task_key = key.clone();
|
||||
let global_registry = std::sync::Arc::clone(&self.global_registry);
|
||||
let _abort_handle = self.tasks.spawn(async move {
|
||||
return fetch_hydration(pool, role, expected_network, task_key, commitment).await;
|
||||
return fetch_hydration_shared(global_registry, pool, role, expected_network, task_key, commitment).await;
|
||||
});
|
||||
}
|
||||
return std::result::Result::Ok(());
|
||||
@@ -3538,6 +3614,53 @@ async fn fetch_hydration(
|
||||
return std::result::Result::Ok(RawTransactionIngestHydrationFetch { key, observed });
|
||||
}
|
||||
|
||||
async fn fetch_hydration_shared(
|
||||
global_registry: std::sync::Arc<RawTransactionIngestGlobalHydrationRegistry>,
|
||||
http_pool: ksp_onchain_transport_lib::HttpTransportPool,
|
||||
hydration_role: ksp_onchain_transport_lib::HttpRoleName,
|
||||
expected_network: ksp_store_lib::RawNetworkId,
|
||||
key: RawTransactionIngestHydrationKey,
|
||||
commitment: ksp_onchain_transport_lib::SolanaCommitment,
|
||||
) -> ksp_core_lib::Result<RawTransactionIngestHydrationFetch> {
|
||||
let (leader, mut receiver) = match global_registry.subscribe_or_lead(&key) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
if leader {
|
||||
let fetched = fetch_hydration(http_pool, hydration_role, expected_network, key.clone(), commitment).await;
|
||||
return match fetched {
|
||||
std::result::Result::Ok(value) => {
|
||||
global_registry.publish_and_remove(&key, RawTransactionIngestSharedHydrationResult::Available(value.observed.clone()));
|
||||
std::result::Result::Ok(value)
|
||||
},
|
||||
std::result::Result::Err(error) => {
|
||||
global_registry.publish_and_remove(&key, RawTransactionIngestSharedHydrationResult::Failed(error.code()));
|
||||
std::result::Result::Err(error)
|
||||
},
|
||||
};
|
||||
}
|
||||
loop {
|
||||
if let std::option::Option::Some(result) = receiver.borrow().clone() {
|
||||
return shared_hydration_result(key, result);
|
||||
}
|
||||
if receiver.changed().await.is_err() {
|
||||
return std::result::Result::Err(crate::runtime_error("source.global_hydration_channel_closed"));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn shared_hydration_result(
|
||||
key: RawTransactionIngestHydrationKey,
|
||||
result: RawTransactionIngestSharedHydrationResult,
|
||||
) -> ksp_core_lib::Result<RawTransactionIngestHydrationFetch> {
|
||||
return match result {
|
||||
RawTransactionIngestSharedHydrationResult::Available(observed) => std::result::Result::Ok(RawTransactionIngestHydrationFetch { key, observed }),
|
||||
RawTransactionIngestSharedHydrationResult::Failed(code) => {
|
||||
std::result::Result::Err(ksp_core_lib::Error::new(code, "RAW transaction ingest shared hydration failed"))
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
fn finalize_hydration(
|
||||
hydration: &RawTransactionIngestHydrationContext,
|
||||
settings: &crate::RawTransactionIngestSettings,
|
||||
|
||||
Reference in New Issue
Block a user