From 1f2720463b16db9bdcd63319c0a1cfdbce66e783 Mon Sep 17 00:00:00 2001 From: SinuS Von SifriduS Date: Thu, 10 Sep 2026 17:52:35 +0200 Subject: [PATCH] v0.3.13-pre.008 --- Cargo.toml | 4 +- .../README.md | 8 +- .../USAGE.md | 4 +- .../src/lib.rs | 6 +- .../src/persistence.rs | 131 +++++++++++- .../src/runtime.rs | 22 +- .../src/runtime_resources.rs | 141 ++++++++++++- .../tests/dependency_boundary.rs | 38 +++- .../tests/hardening.rs | 37 +++- .../tests/public_api.rs | 17 +- .../tests/release_completeness.rs | 5 +- .../unit_tests/admission.rs | 51 ++++- .../unit_tests/persistence.rs | 110 +++++++++- .../unit_tests/runtime.rs | 32 ++- .../unit_tests/runtime_resources.rs | 44 +++- deltas/0.3.13/pre.008.md | 198 ++++++++++++++++++ ...3_13_MULTI_SOURCE_LIVE_CONVERGENCE_PLAN.md | 39 +++- ...0-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE.md | 48 ++++- 18 files changed, 899 insertions(+), 36 deletions(-) create mode 100644 deltas/0.3.13/pre.008.md diff --git a/Cargo.toml b/Cargo.toml index e7639ca..92d9166 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 552 +# version: 553 [workspace] resolver = "3" members = ["crates/ksp-app-backfill-desk", "crates/ksp-app-config-desk", "crates/ksp-app-solprices-desk", "crates/ksp-app-store-desk", "crates/ksp-app-wallet-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-interface-lib", "crates/ksp-job-api", "crates/ksp-job-backfill-lib", "crates/ksp-logging-lib", "crates/ksp-offchain-transport-lib", "crates/ksp-onchain-transport-lib", "crates/ksp-program-api", "crates/ksp-raw-transaction-lib", "crates/ksp-store-api", "crates/ksp-store-lib", "crates/ksp-store-postgres-lib", "crates/ksp-wallet-lib", "crates/ksp-worker-api", "crates/ksp-worker-raw-transaction-ingest-lib"] [workspace.package] -version = "0.3.13-pre.7" +version = "0.3.13-pre.8" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/README.md b/crates/ksp-worker-raw-transaction-ingest-lib/README.md index 9b6403c..8916986 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/README.md +++ b/crates/ksp-worker-raw-transaction-ingest-lib/README.md @@ -1,5 +1,5 @@ - + # ksp-worker-raw-transaction-ingest-lib @@ -86,7 +86,7 @@ Le runtime-resource aggregate public accepte une collection validée de 1 à 32 Le supervisor possède toutes les tâches source. Une source qui échoue est terminale pour le Worker et déclenche l'arrêt coopératif puis le join des autres sources, car cette release ne suppose aucune équivalence de coverage. Une fermeture propre d'une source alors que le Worker n'est pas en arrêt est également traitée comme une perte de source configurée et devient terminale. -Chaque source conserve provisoirement ses mécanismes de production/hydration existants. Un inventaire privé `source_key -> latest processing/source state`, borné à 32 entrées, agrège la projection run-local. La frontier agrégée reste conservative : elle n'expose un `processing_frontier_slot` que lorsque toutes les sources en possèdent un, choisit le minimum des frontiers connus et le plus ancien pending. La convergence/coalescence cross-source des mêmes transactions reste une tranche séparée. +Un inventaire privé `source_key -> latest processing/source state`, borné à 32 entrées, agrège la projection run-local. La frontier agrégée reste conservative : elle n'expose un `processing_frontier_slot` que lorsque toutes les sources en possèdent un, choisit le minimum des frontiers connus et le plus ancien pending. Les sources reference-bearing partagent en plus un registre global d'hydration borné : une même clé `(network, signature, commitment)` ne déclenche qu'un leader HTTP, puis chaque signal source conserve sa propre provenance lors de la finalisation. ## Contrat de source Standard Logs + HTTP @@ -171,14 +171,14 @@ Le shutdown est borné par `shutdown_drain_timeout`. Le supervisor multi-source La queue centrale est un `tokio::sync::mpsc` privé borné par `admission_queue_capacity`. Les sources internes subissent la backpressure ; aucune queue non bornée ni silent drop n'est autorisé. -Le Worker possède un coordinateur d'hydration source-neutral borné, réutilisé par Yellowstone, Standard Logs et Helius Transaction. Standard Block n'entre pas dans ce coordinateur lorsqu'une transaction est direct-qualified : +Le Worker possède des coordinateurs source-neutral Yellowstone, Standard Logs et Helius Transaction branchés sur un registre global d'hydration partagé. Standard Block et HTTP Block Polling n'entrent pas dans ce registre lorsqu'une transaction est direct-qualified : ```text in-flight hydration <= persistence_concurrency pending source signals <= admission_queue_capacity ``` -Les signaux partageant le même `(network, signature, commitment)` sont coalescés avant le fan-out HTTP. Les provenances utiles restent néanmoins conservées pour les ingress produits après hydration. +Les signaux partageant le même `(network, signature, commitment)` sont coalescés cross-source avant le fan-out HTTP. Après canonicalisation, une cache run-local bornée sérialise les acquisitions de même `(network, signature)` : la première passe par l'écriture atomique entity + observation, les suivantes de contenu identique ajoutent uniquement leur observation déterministe. Un hash canonique divergent devient un content conflict terminal. Les retries/reroutages HTTP appartiennent à `ksp-onchain-transport-lib`. Le Worker ne possède pas une seconde boucle de retry autour de `getTransaction`. diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md b/crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md index 47d3336..7d86a32 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md +++ b/crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md @@ -1,5 +1,5 @@ - + # Utilisation de ksp-worker-raw-transaction-ingest-lib @@ -253,7 +253,7 @@ Yellowstone Slot -> continuity-only Yellowstone Account/Ping/Pong/Entry -> sans RAW Transaction dans cette verticale ``` -Les signaux de même `(network, signature, commitment)` sont coalescés avant l'hydration HTTP. Le Worker ne possède pas une boucle de retry HTTP : reroutage/retry/backoff restent dans `ksp-onchain-transport-lib`. +Les signaux de même `(network, signature, commitment)` sont coalescés globalement avant l'hydration HTTP. Si plusieurs sources produisent ensuite la même transaction canonique, le Worker conserve une seule entité RAW et enregistre séparément les observations déterministes propres à chaque source. Le Worker ne possède pas une boucle de retry HTTP : reroutage/retry/backoff restent dans `ksp-onchain-transport-lib`. ## Observer le snapshot concret diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs index a082ec9..8e6d4c8 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs @@ -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. diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs index 3c71809..1d3dc4f 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs @@ -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>; + + /// 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>; } 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> { 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> { + 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>>>, + >, +} + +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>>> { + 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

( + port: &P, + acquisition: ksp_raw_transaction_lib::RawTransactionAcquisition, + convergence: &crate::RawTransactionIngestPersistenceConvergence, +) -> ksp_core_lib::Result +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 { 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 { + 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; diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs index d1e4a31..d38a07c 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs -// version: 14 +// version: 15 type PersistencePort = std::sync::Arc; type PersistenceTasks = tokio::task::JoinSet>; @@ -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, port: &std::option::Option, snapshots: &mut crate::RawTransactionIngestSnapshotPublisher, ) -> std::option::Option { @@ -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, port: &std::option::Option, snapshots: &mut crate::RawTransactionIngestSnapshotPublisher, ) -> std::option::Option { 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( } 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( &mut children, &mut admission, &mut persistence, + &persistence_convergence, &port, &mut snapshots, &mut processing_frontier_receiver, @@ -399,7 +405,8 @@ async fn run_supervisor( 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( fn spawn_persistence( persistence: &mut PersistenceTasks, + persistence_convergence: &std::sync::Arc, port: &std::option::Option, 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, port: &std::option::Option, snapshots: &mut crate::RawTransactionIngestSnapshotPublisher, processing_frontier_receiver: &mut std::option::Option, @@ -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(), diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs index 86587ef..ced1ec9 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs @@ -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, admission_sender: tokio::sync::mpsc::Sender, + global_hydration_registry: std::sync::Arc, 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, admission_sender: tokio::sync::mpsc::Sender, processing_frontier_sender: tokio::sync::watch::Sender, + global_hydration_registry: std::sync::Arc, ) -> 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, admission_sender: tokio::sync::mpsc::Sender, processing_frontier_sender: tokio::sync::watch::Sender, + global_hydration_registry: std::sync::Arc, ) -> 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, admission_sender: tokio::sync::mpsc::Sender, processing_frontier_sender: tokio::sync::watch::Sender, + global_hydration_registry: std::sync::Arc, ) -> 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::>(); 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>, + >, + >, +} + +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>)> { + 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>; struct RawTransactionIngestHydrationCoordinator { + global_registry: std::sync::Arc, 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, + ) -> 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, + 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 { + 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 { + 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, diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs b/crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs index d256c53..c7fa7ad 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs -// version: 23 +// version: 24 //! Dependency firewall canaries for the RAW transaction ingest Worker foundation. @@ -568,3 +568,39 @@ fn v0_3_12_pre_008_worker_observes_transport_reconnect_replay_and_faults_only_on } return; } + +#[test] +fn v0_3_13_pre_008_cross_source_convergence_reuses_store_and_transport_facades_only() { + let root = include_str!("../src/lib.rs"); + let persistence = include_str!("../src/persistence.rs"); + let runtime = include_str!("../src/runtime.rs"); + let resources = include_str!("../src/runtime_resources.rs"); + for required in [ + "RawTransactionIngestPersistenceConvergence", + "persist_raw_transaction_ingest_converged_acquisition", + "RawTransactionIngestPersistencePort", + "record_raw_transaction_observation", + "RawTransactionObservationWrite", + "payload().content_hash()", + "content_conflict_error()", + ] { + assert!(persistence.contains(required) || root.contains(required), "required pre.008 persistence convergence contract missing: {required}"); + } + for required in [ + "RawTransactionIngestGlobalHydrationRegistry", + "fetch_hydration_shared", + "subscribe_or_lead", + "get_transaction_observed", + "settings.admission_queue_capacity()", + ] { + assert!(resources.contains(required), "required pre.008 hydration convergence contract missing: {required}"); + } + assert!(runtime.contains("settings.admission_queue_capacity().max(settings.persistence_concurrency())")); + for forbidden in ["ksp_store_postgres_lib::", "reqwest::", "tonic::", "yellowstone_grpc_proto::", "ksp_job_backfill_lib::", "unbounded_channel"] { + assert!( + !persistence.contains(forbidden) && !runtime.contains(forbidden) && !resources.contains(forbidden), + "pre.008 convergence crossed a dependency/boundedness boundary: {forbidden}" + ); + } + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs b/crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs index a6fc04c..40468dc 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs -// version: 17 +// version: 18 //! External public, security, redaction and release-boundary hardening canaries for `pre.010`. @@ -751,3 +751,38 @@ fn v0_3_13_pre_007_multi_source_supervisor_is_fail_closed_joined_and_does_not_pu } return; } + +#[test] +fn v0_3_13_pre_008_cross_source_convergence_is_bounded_conflict_checked_and_private() { + let root = include_str!("../src/lib.rs"); + let persistence = include_str!("../src/persistence.rs"); + let resources = include_str!("../src/runtime_resources.rs"); + for required in [ + "max_entries: usize", + "entries.len() >= self.max_entries", + "std::sync::Arc::strong_count(entry) == 1", + "known_hash_value != content_hash", + "record_observation(observation)", + "persistence.additional_observation_not_recorded", + ] { + assert!(persistence.contains(required), "required pre.008 persistence hardening guard missing: {required}"); + } + for required in [ + "max_pending: usize", + "pending.len() >= self.max_pending", + "source.global_hydration_pending_saturated", + "tokio::sync::watch::channel", + "publish_and_remove", + ] { + assert!(resources.contains(required), "required pre.008 hydration hardening guard missing: {required}"); + } + for forbidden in [ + "pub struct RawTransactionIngestPersistenceConvergence", + "pub struct RawTransactionIngestGlobalHydrationRegistry", + "pub fn persist_raw_transaction_ingest_converged_acquisition", + ] { + assert!(!root.contains(forbidden), "pre.008 private convergence implementation leaked publicly: {forbidden}"); + } + assert!(!resources.contains("unbounded_channel")); + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/tests/public_api.rs b/crates/ksp-worker-raw-transaction-ingest-lib/tests/public_api.rs index 7081d80..310bbc8 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/tests/public_api.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/tests/public_api.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/tests/public_api.rs -// version: 17 +// version: 18 //! External public-surface proofs for the RAW transaction ingest Worker foundation. @@ -336,3 +336,18 @@ fn v0_3_13_pre_007_source_inventory_and_logical_keys_remain_private() { } return; } + +#[test] +fn v0_3_13_pre_008_convergence_cache_registry_and_source_keys_remain_private() { + let root = include_str!("../src/lib.rs"); + for forbidden in [ + "pub use self::persistence::RawTransactionIngestPersistenceConvergence", + "pub use self::persistence::persist_raw_transaction_ingest_converged_acquisition", + "RawTransactionIngestGlobalHydrationRegistry", + "RawTransactionIngestHydrationKey", + "source_key", + ] { + assert!(!root.contains(forbidden), "pre.008 convergence implementation leaked through public crate root: {forbidden}"); + } + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs b/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs index 8dac106..464e919 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs -// version: 15 +// version: 16 //! Release-completeness canaries through the `pre.010` public/release/security hardening tranche. @@ -128,6 +128,7 @@ fn pre_010_external_hardening_suite_is_present_and_scoped() { "v0_3_13_pre_005_helius_transaction_redaction_full_reference_and_tier_neutrality_are_explicit", "v0_3_13_pre_006_http_block_polling_is_bounded_run_local_stop_preemptible_and_redacted", "v0_3_13_pre_007_multi_source_supervisor_is_fail_closed_joined_and_does_not_publish_source_identity", + "v0_3_13_pre_008_cross_source_convergence_is_bounded_conflict_checked_and_private", ] { assert!(hardening.contains(required), "required pre.010 hardening canary missing: {required}"); } @@ -146,6 +147,7 @@ fn pre_010_external_hardening_suite_is_present_and_scoped() { assert!(dependency_boundary.contains("v0_3_13_pre_005_helius_transaction_source_reuses_transport_facade_and_common_hydration_only")); assert!(dependency_boundary.contains("v0_3_13_pre_006_http_block_polling_reuses_transport_facades_without_becoming_backfill")); assert!(dependency_boundary.contains("v0_3_13_pre_007_multi_source_supervisor_and_inventory_remain_private_bounded_and_source_neutral")); + assert!(dependency_boundary.contains("v0_3_13_pre_008_cross_source_convergence_reuses_store_and_transport_facades_only")); let public_api = include_str!("public_api.rs"); assert!(public_api.contains("pre_003_kind_code_and_settings_are_consumable_from_crate_root")); assert!(public_api.contains("pre_004_start_handle_and_terminal_future_are_consumable_without_public_join_handle")); @@ -157,6 +159,7 @@ fn pre_010_external_hardening_suite_is_present_and_scoped() { assert!(public_api.contains("v0_3_13_pre_004_standard_block_runtime_resource_surface_is_typed_and_transport_owned")); assert!(public_api.contains("v0_3_13_pre_006_http_block_polling_runtime_resource_surface_is_bounded_and_transport_owned")); assert!(public_api.contains("v0_3_13_pre_007_source_inventory_and_logical_keys_remain_private")); + assert!(public_api.contains("v0_3_13_pre_008_convergence_cache_registry_and_source_keys_remain_private")); assert!(public_api.contains("v0_3_12_pre_007_processing_frontier_snapshot_getters_are_public_and_processing_only")); assert!(public_api.contains("pre_008_snapshot_surface_and_common_projection_are_public_and_stable")); return; diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs index efa60b2..9ea970d 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs -// version: 4 +// version: 5 fn material( network: &str, @@ -262,3 +262,52 @@ async fn pre_009_full_queue_dequeue_marks_source_neutral_backpressure_observatio assert!(!admission.take_backpressure_wait_observed()); return; } + +#[tokio::test(flavor = "current_thread")] +async fn v0_3_13_pre_008_distinct_source_keys_produce_distinct_idempotent_observation_keys() { + let (network, first_material) = match material("mainnet", 51, "AQID") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let (_, second_material) = match material("mainnet", 51, "AQID") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let first_provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let second_provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let (mut admission, sender) = crate::RawTransactionAdmission::new(2); + let first_ingress = crate::RawTransactionIngress { + material: first_material, + network: network.clone(), + provenance: first_provenance, + source_key: [11_u8; 32], + }; + let second_ingress = crate::RawTransactionIngress { + material: second_material, + network: network.clone(), + provenance: second_provenance, + source_key: [12_u8; 32], + }; + if sender.send(first_ingress).await.is_err() || sender.send(second_ingress).await.is_err() { + return; + } + std::mem::drop(sender); + let first = match admission.receive(&network).await { + std::result::Result::Ok(std::option::Option::Some(value)) => value, + _ => return, + }; + let second = match admission.receive(&network).await { + std::result::Result::Ok(std::option::Option::Some(value)) => value, + _ => return, + }; + assert_eq!(first.transaction().reference(), second.transaction().reference()); + assert_eq!(first.transaction().payload().content_hash(), second.transaction().payload().content_hash()); + assert_ne!(first.observation().observation_key(), second.observation().observation_key()); + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/persistence.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/persistence.rs index 0dbb468..0510c03 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/persistence.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/persistence.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/persistence.rs -// version: 1 +// version: 2 #[derive(Clone, Copy)] enum PortResponse { @@ -12,6 +12,8 @@ struct FakePersistencePort { calls: std::sync::atomic::AtomicUsize, network: ksp_store_lib::RawNetworkId, normal_mode_seen: std::sync::atomic::AtomicBool, + observation_calls: std::sync::atomic::AtomicUsize, + observation_response: std::option::Option, response: PortResponse, } @@ -21,9 +23,16 @@ impl FakePersistencePort { calls: std::sync::atomic::AtomicUsize::new(0), network, normal_mode_seen: std::sync::atomic::AtomicBool::new(false), + observation_calls: std::sync::atomic::AtomicUsize::new(0), + observation_response: std::option::Option::None, response, }; } + + fn with_observation_response(mut self, observation_response: ksp_store_lib::RawObservationWriteOutcome) -> Self { + self.observation_response = std::option::Option::Some(observation_response); + return self; + } } impl crate::RawTransactionIngestPersistencePort for FakePersistencePort { @@ -51,9 +60,39 @@ impl crate::RawTransactionIngestPersistencePort for FakePersistencePort { }; }); } + + fn record_observation<'a>( + &'a self, + _observation: ksp_store_lib::RawTransactionObservation, + ) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result> { + self.observation_calls.fetch_add(1, std::sync::atomic::Ordering::AcqRel); + let observation_response = self.observation_response; + let response = self.response; + return std::boxed::Box::pin(async move { + if let std::option::Option::Some(outcome) = observation_response { + return std::result::Result::Ok(outcome); + } + return match response { + PortResponse::Outcome(_, observation) => std::result::Result::Ok(observation), + PortResponse::Conflict => std::result::Result::Err(ksp_core_lib::Error::new(ksp_store_lib::ERROR_CODE_RAW_CONFLICT, "synthetic conflict")), + PortResponse::StoreFailure => std::result::Result::Err(ksp_core_lib::Error::new( + ksp_core_lib::ErrorCode::new("synthetic_store", "write_failed"), + "synthetic remote failure text that must not escape", + )), + }; + }); + } } fn acquisition(network_name: &str, signature_byte: u8) -> std::option::Option { + return acquisition_with_observation_key(network_name, signature_byte, signature_byte); +} + +fn acquisition_with_observation_key( + network_name: &str, + signature_byte: u8, + observation_byte: u8, +) -> std::option::Option { let network = match ksp_store_lib::RawNetworkId::new(network_name) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, @@ -92,7 +131,7 @@ fn acquisition(network_name: &str, signature_byte: u8) -> std::option::Option value, + std::option::Option::None => return, + }; + let port = FakePersistencePort::new( + network, + PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::Inserted), + ); + let convergence = crate::RawTransactionIngestPersistenceConvergence::new(8); + let first = match acquisition_with_observation_key("mainnet", 41, 1) { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let second = match acquisition_with_observation_key("mainnet", 41, 2) { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let first_outcome = crate::persist_raw_transaction_ingest_converged_acquisition(&port, first, &convergence).await; + let second_outcome = crate::persist_raw_transaction_ingest_converged_acquisition(&port, second, &convergence).await; + assert!(first_outcome.is_ok()); + assert!(second_outcome.is_ok()); + assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 1); + assert_eq!(port.observation_calls.load(std::sync::atomic::Ordering::Acquire), 1); + let second_outcome = match second_outcome { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert_eq!(second_outcome.entity(), crate::RawTransactionIngestEntityPersistence::AlreadyPresent); + assert_eq!(second_outcome.observation(), crate::RawTransactionIngestObservationPersistence::Inserted); + return; +} + +#[tokio::test(flavor = "current_thread")] +async fn v0_3_13_pre_008_replayed_same_observation_is_idempotent_without_raw_resubmission() { + let network = match network("mainnet") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let port = FakePersistencePort::new( + network, + PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::Inserted), + ) + .with_observation_response(ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent); + let convergence = crate::RawTransactionIngestPersistenceConvergence::new(8); + let first = match acquisition_with_observation_key("mainnet", 42, 3) { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let second = match acquisition_with_observation_key("mainnet", 42, 3) { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let first_result = crate::persist_raw_transaction_ingest_converged_acquisition(&port, first, &convergence).await; + assert!(first_result.is_ok()); + let second_result = crate::persist_raw_transaction_ingest_converged_acquisition(&port, second, &convergence).await; + let second_outcome = match second_result { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 1); + assert_eq!(port.observation_calls.load(std::sync::atomic::Ordering::Acquire), 1); + assert_eq!(second_outcome.observation(), crate::RawTransactionIngestObservationPersistence::AlreadyPresent); + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs index 7d0ab92..6cc99af 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs -// version: 9 +// version: 10 struct ActiveTaskGuard { active: std::sync::Arc, @@ -348,6 +348,36 @@ impl crate::RawTransactionIngestPersistencePort for RuntimePersistencePort { )); }); } + + fn record_observation<'a>( + &'a self, + _observation: ksp_store_lib::RawTransactionObservation, + ) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result> { + let response = self.response; + return std::boxed::Box::pin(async move { + if matches!(response, RuntimePortResponse::StoreFailureBlocked | RuntimePortResponse::SuccessBlocked) { + let current = self.active.fetch_add(1, std::sync::atomic::Ordering::AcqRel) + 1; + let _active_guard = RuntimePortActiveGuard { active: &self.active }; + self.update_max_active(current); + while !self.released.load(std::sync::atomic::Ordering::Acquire) { + self.notify.notified().await; + } + if matches!(response, RuntimePortResponse::SuccessBlocked) { + self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); + return std::result::Result::Ok(ksp_store_lib::RawObservationWriteOutcome::Inserted); + } + } + if matches!(response, RuntimePortResponse::Conflict) { + self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); + return std::result::Result::Err(ksp_core_lib::Error::new(ksp_store_lib::ERROR_CODE_RAW_CONFLICT, "synthetic conflict")); + } + self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); + return std::result::Result::Err(ksp_core_lib::Error::new( + ksp_core_lib::ErrorCode::new("synthetic_store", "write_failed"), + "synthetic store failure", + )); + }); + } } fn runtime_ingress(network: &ksp_store_lib::RawNetworkId, signature_byte: u8) -> std::option::Option { diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs index cccc5f9..999bdb1 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs -// version: 19 +// version: 20 fn grpc_endpoint(cluster: &str) -> std::option::Option { return grpc_endpoint_with_identity(cluster, "yellowstone-fixture", "fixture-provider"); @@ -2992,3 +2992,45 @@ fn v0_3_13_pre_006_runtime_resources_reject_duplicate_http_polling_identity_even assert_eq!(resources.source_count(), 1); return; } + +#[test] +fn v0_3_13_pre_008_global_hydration_registry_coalesces_cross_source_key_to_one_leader() { + let settings = match pre_004_settings() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let network = match ksp_store_lib::RawNetworkId::new("devnet") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let key = super::RawTransactionIngestHydrationKey { + commitment: ksp_onchain_transport_lib::SolanaCommitment::Confirmed.as_str(), + network, + signature: ksp_store_lib::RawTransactionSignature::new([61_u8; 64]), + }; + let registry = super::RawTransactionIngestGlobalHydrationRegistry::new(&settings); + let (first_leader, _first_receiver) = match registry.subscribe_or_lead(&key) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let (second_leader, second_receiver) = match registry.subscribe_or_lead(&key) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert!(first_leader); + assert!(!second_leader); + registry.publish_and_remove(&key, super::RawTransactionIngestSharedHydrationResult::Failed(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_SOURCE_FAILED)); + let shared = second_receiver.borrow().clone(); + match shared { + std::option::Option::Some(super::RawTransactionIngestSharedHydrationResult::Failed(code)) => { + assert_eq!(code, crate::ERROR_CODE_RAW_TRANSACTION_INGEST_SOURCE_FAILED); + }, + _ => return, + } + let (third_leader, _third_receiver) = match registry.subscribe_or_lead(&key) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert!(third_leader); + return; +} diff --git a/deltas/0.3.13/pre.008.md b/deltas/0.3.13/pre.008.md new file mode 100644 index 0000000..cdba864 --- /dev/null +++ b/deltas/0.3.13/pre.008.md @@ -0,0 +1,198 @@ + + + +# Delta 0.3.13-pre.008 — convergence cross-source et observations multiples + +## Base requise + +```text +livraison précédente : 0.3.13-pre.007 +Cargo base : 0.3.13-pre.7 +delta base : deltas/0.3.13/pre.007.md +archive delta base : ksp-general-0.3.13-pre.007.zip +``` + +Le gate opérateur communiqué le 10 septembre 2026 est vert pour toutes les commandes exécutées : `cargo fmt`, `cargo fmt --check`, audits Rust/Markdown, `cargo check --workspace`, Clippy strict et toutes les suites `ksp-worker-raw-transaction-ingest-lib` (`90` unit, `4` cross-layer, `14` dependency-boundary, `21` hardening, `15` public-api, `4` release-completeness, `0` doc-test). Le log reçu ne contient pas les sorties séparées `cargo test -p ksp-onchain-transport-lib` ni les trois commandes `cargo tree`; elles ne sont donc pas déclarées PASS ici. + +## Objectif + +Fermer la convergence des sources live simultanées sans dupliquer l'identité `RawTransaction` : + +```text +plusieurs sources reference-bearing + -> une hydration globale par (network, signature, commitment) + -> acquisitions source-distinctes + -> canonicalisation RAW commune + -> convergence par (network, signature) + content hash + -> une entity RAW + -> plusieurs observations déterministes +``` + +## Version + +```text +livraison : 0.3.13-pre.008 +workspace.package.version : 0.3.13-pre.8 +archive : ksp-general-0.3.13-pre.008.zip +``` + +`Cargo.toml` passe de la version de fichier `552` à `553`. + +## Hydration globale cross-source + +`RawTransactionIngestRuntimeResources::run_live_sources` crée un `RawTransactionIngestGlobalHydrationRegistry` privé partagé par Yellowstone, Standard Logs et Helius Transaction. Standard Block et HTTP Block Polling restent RAW-direct et n'entrent pas dans ce registre. + +La clé reste exactement : + +```text +(network, signature, commitment) +``` + +Le premier demandeur devient leader et exécute le `getTransaction observed` existant. Les suivants s'abonnent au même résultat via un `watch` privé. La registry est bornée par `admission_queue_capacity`, publie soit le résultat observed soit un `ErrorCode` sûr, puis retire la clé. + +Les signaux source restent distincts : le même résultat d'hydration est finalisé séparément avec le contexte/provenance de chaque signal. + +## Convergence persistence + +Le runtime possède un `RawTransactionIngestPersistenceConvergence` privé et run-local. Sa capacité productive est : + +```text +max(admission_queue_capacity, persistence_concurrency) +``` + +Chaque identité `(network, signature)` possède un verrou async local et mémorise le `RawContentHash` canonique après une première persistence réussie non purgée. + +```text +première acquisition + -> persist_raw_transaction_acquisition + -> Store atomique entity + observation + +même identity + même content hash + -> record_raw_transaction_observation uniquement + +même identity + hash divergent + -> content conflict terminal +``` + +Le cache évince uniquement une entrée inactive lorsqu'il atteint sa borne. Le runtime dimensionne cette borne pour couvrir toutes les tâches persistence simultanées valides. + +## Observations multiples + +L'`observation_key` reste dérivée de la référence canonique et du `source_key` privé. Deux sources distinctes observant le même RAW produisent donc une seule référence/entity mais des observations déterministes distinctes. + +Un replay de la même observation retourne `AlreadyPresent` sans réécriture de l'entité. Une observation supplémentaire `NotRecorded` est incohérente pour une entité déjà durable et devient une faute runtime sûre. + +## Preuves ajoutées + +Tests unitaires déterministes : + +```text +deux source_key -> même référence/hash canonique, observation_key distinctes +registry globale -> un leader, un follower, publication partagée puis retrait +première identity -> une persistence atomique +seconde observation distincte -> record_observation seul +replay même observation -> AlreadyPresent sans nouvelle persistence entity +``` + +Canaris externes : + +```text +dependency_boundary : Store/Transport facades uniquement, pas de backend/client direct +hardening : caches bornés, hash conflict-checked, convergence privée, aucun canal unbounded +public_api : registry/cache/hydration key/source_key non exposés +release_completeness : présence obligatoire des trois canaris pre.008 +``` + +## Frontière de tranche + +`pre.008` ne ferme pas encore : + +```text +fairness adversariale sous duplicate storm prolongé +starvation entre sources +matrice complète de disagreement provider/source +stress des bornes globales +nouveaux compteurs/health publics multi-source +``` + +Ces points restent `pre.009` et `pre.010` selon le plan. + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-worker-raw-transaction-ingest-lib/README.md +crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md +crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs +crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs +crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs +crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs +crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs +crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs +crates/ksp-worker-raw-transaction-ingest-lib/tests/public_api.rs +crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs +crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs +crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/persistence.rs +crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs +crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs +docs/plans/034-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE_PLAN.md +docs/validation/030-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE.md +``` + +## Fichiers ajoutés + +```text +deltas/0.3.13/pre.008.md +``` + +## Fichiers supprimés + +```text +aucun +``` + +## Validation dans l'environnement d'assemblage + +Le toolchain Rust/Cargo/Rustfmt n'est pas disponible dans l'environnement d'assemblage. Les commandes Cargo post-`pre.008` restent donc `NON EXÉCUTÉ LOCAL`. Les audits statiques KSP, le scan normatif, le contrôle exhaustif du diff, des headers et du ZIP sont exécutés avant livraison. + +Résultats statiques finaux : + +```text +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean +Markdown table audit: clean (340 tables, 842 files) +Normative rule definitions: 489 +Unique normative IDs: 489 +Duplicates: 0 +``` + +Contrôle ciblé : + +```text +17 fichiers existants modifiés +1 fichier ajouté +0 suppression +17/17 headers existants : +1 +nouvelle dépendance : aucune +backend Store physique direct : absent +client HTTP/gRPC direct : absent +unbounded_channel : absent +convergence/registry publiques : absentes +``` + +## Gate opérateur requis avant pre.009 + +```bash +cargo fmt --all +cargo fmt --all -- --check +python3 scripts/audit_rust_workspace_rules.py +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas +cargo check --workspace +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo test -p ksp-onchain-transport-lib +cargo test -p ksp-worker-raw-transaction-ingest-lib +cargo tree -p ksp-worker-raw-transaction-ingest-lib --edges normal +cargo tree -p ksp-worker-raw-transaction-ingest-lib -e features +cargo tree --duplicates +``` diff --git a/docs/plans/034-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE_PLAN.md b/docs/plans/034-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE_PLAN.md index 160a8d2..e80a72f 100644 --- a/docs/plans/034-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE_PLAN.md +++ b/docs/plans/034-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE_PLAN.md @@ -1,5 +1,5 @@ - + # Plan v0.3.13 — WS standard / Helius / HTTP live + convergence multi-source RawTransaction @@ -1246,3 +1246,40 @@ cargo tree -p ksp-worker-raw-transaction-ingest-lib -e features cargo tree --duplicates ``` +## 62. Gate opérateur pre.007 reçu avant pre.008 + +Le gate communiqué le 10 septembre 2026 est vert pour toutes les commandes effectivement exécutées : fmt, audits Rust/Markdown, `cargo check --workspace`, Clippy strict et toutes les suites Worker (`90` unit, `4` cross-layer, `14` dependency-boundary, `21` hardening, `15` public-api, `4` release-completeness, `0` doc-test). Les sorties Transport ciblées et `cargo tree` ne figurent pas dans le log reçu et ne sont pas déclarées PASS. + +## 63. Implémentation pre.008 — convergence cross-source + +Les trois familles reference-bearing Yellowstone, Standard Logs et Helius Transaction partagent un `RawTransactionIngestGlobalHydrationRegistry` privé au run. La clé reste exactement `(network, signature, commitment)`. La première source devient leader et exécute `getTransaction observed`; les suivantes attendent le même résultat via un canal `watch` privé. Le registre est borné par `admission_queue_capacity` et retiré après publication du résultat. + +La canonicalisation reste inchangée et source-neutral. La convergence persistence utilise `(network, signature)` puis compare le `RawContentHash` canonique. Une divergence de hash devient le content conflict terminal existant. + +## 64. Observations multiples et persistence + +La première acquisition d'une identité canonique passe par l'écriture Store atomique existante `RawTransaction + RawTransactionObservation`. Tant que le hash reste identique, les acquisitions suivantes de la même identité n'écrivent plus l'entité : elles appellent uniquement `record_raw_transaction_observation` via `ksp-store-lib`. + +Les `RawObservationKey` restent déterministes et source-dépendantes. Deux sources distinctes peuvent donc produire deux observations durables pour une seule entité RAW, tandis qu'un replay de la même observation reste idempotent (`AlreadyPresent`). `NotRecorded` sur une observation supplémentaire est considéré incohérent et devient une faute runtime sûre. + +La cache persistence est privée, run-local et bornée à `max(admission_queue_capacity, persistence_concurrency)` dans le runtime productif afin de couvrir toutes les tâches persistence simultanées sans bypass de convergence. + +## 65. Frontière de tranche pre.008 + +`pre.008` ferme la coalescence cross-source, l'unique hydration globale par clé et les observations Store multiples. Elle ne ferme pas encore les scénarios de disagreement/fairness sous charge hostile, les duplicate storms prolongés, la starvation ou les bornes adversariales détaillées : ces preuves appartiennent à `pre.009`. + +## 66. Gate opérateur requis avant pre.009 + +```bash +cargo fmt --all +cargo fmt --all -- --check +python3 scripts/audit_rust_workspace_rules.py +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas +cargo check --workspace +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo test -p ksp-onchain-transport-lib +cargo test -p ksp-worker-raw-transaction-ingest-lib +cargo tree -p ksp-worker-raw-transaction-ingest-lib --edges normal +cargo tree -p ksp-worker-raw-transaction-ingest-lib -e features +cargo tree --duplicates +``` diff --git a/docs/validation/030-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE.md b/docs/validation/030-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE.md index 7b6ca30..00b2ed3 100644 --- a/docs/validation/030-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE.md +++ b/docs/validation/030-V0_3_13_MULTI_SOURCE_LIVE_CONVERGENCE.md @@ -1,5 +1,5 @@ - + # Validation v0.3.13 — WS standard / Helius / HTTP live + convergence multi-source @@ -1242,3 +1242,49 @@ cargo tree -p ksp-worker-raw-transaction-ingest-lib -e features cargo tree --duplicates ``` +## 63. Gate opérateur pre.007 reçu + +Preuve communiquée le 10 septembre 2026 : fmt PASS, audits Rust clean/export completeness 0, Markdown clean `340/841`, `cargo check --workspace` PASS, Clippy strict PASS, Worker `90/90` unit, `4/4` cross-layer, `14/14` dependency-boundary, `21/21` hardening, `15/15` public-api, `4/4` release-completeness et `0` doc-test. Le log fourni ne contient pas les sorties Transport ciblées ni les trois `cargo tree`. + +## 64. Surface pre.008 matérialisée + +```text +RawTransactionIngestGlobalHydrationRegistry privé +coalescence globale (network, signature, commitment) +leader unique getTransaction observed +résultat partagé par watch borné/retiré +RawTransactionIngestPersistenceConvergence privé +clé canonique (network, signature) +hash canonique conflict-checked +première persistence atomique entity + observation +observations suivantes via record_raw_transaction_observation +RawObservationKey source-distincte et idempotente +``` + +## 65. Invariants pre.008 + +```text +aucun fanout HTTP pour une même clé d'hydration simultanée +aucune fusion des identités de source dans l'API publique +une identité RAW canonique n'est pas dupliquée par source +hash divergent pour même network/signature => content conflict terminal +observation source différente => observation durable supplémentaire +replay même observation => AlreadyPresent sans réécriture RAW +registre hydration <= admission_queue_capacity +cache persistence runtime <= max(admission_queue_capacity, persistence_concurrency) +aucun backend Store physique ni client HTTP direct dans Worker +``` + +## 66. Validation locale d'assemblage pre.008 + +Le toolchain Cargo/Rustfmt n'est pas disponible dans l'environnement d'assemblage. Les commandes Cargo post-modification restent `NON EXÉCUTÉ LOCAL`; les audits statiques et contrôles de packaging sont consignés dans le delta `pre.008`. + +## 67. Non-claims pre.008 + +```text +pas encore de fairness adversariale complète avant pre.009 +pas encore de canari starvation/duplicate storm cross-source complet avant pre.009 +pas de failover provider implicite +pas de préférence silencieuse entre sources divergentes +pas de nouveaux compteurs/health publics multi-source avant pre.010 +```