v0.3.15-pre.009-fix.001

This commit is contained in:
2026-09-14 14:47:04 +02:00
parent 11105fac28
commit d7c2931c10
10 changed files with 128 additions and 19 deletions

View File

@@ -1,5 +1,5 @@
<!-- file: crates/ksp-worker-raw-transaction-ingest-lib/README.md -->
<!-- version: 13 -->
<!-- version: 14 -->
# ksp-worker-raw-transaction-ingest-lib
@@ -165,7 +165,7 @@ Le Worker s'exécute sur le runtime Tokio courant du caller. Il ne crée pas de
- utiliser la même source via `WorkerSnapshotSource` ;
- attendre le terminal après drain et join des tâches possédées.
Le shutdown est borné par `shutdown_drain_timeout`. Le supervisor multi-source relaie le stop à toutes les sources et les rejoint avant de rendre son résultat au supervisor Worker ; les tâches source, hydration et persistence possédées sont ensuite drainées ou abort+join avant publication terminale. Si la deadline expire, l'abort du wrapper source détruit aussi son `JoinSet` interne et annule ses tâches imbriquées avant le terminal. Une faute déjà observée n'est pas remplacée par un stop concurrent, sauf le `drain_timeout` terminal lorsqu'une récupération bornée dépasse sa deadline. L'abandon terminal d'une hydration retire son pending run-local sans le convertir artificiellement en travail `settled`.
Le shutdown est borné par `shutdown_drain_timeout` (10 s par défaut). Le supervisor multi-source relaie le stop à toutes les sources et les rejoint avant de rendre son résultat au supervisor Worker ; les tâches source, hydration et persistence possédées sont ensuite drainées ou abort+join avant publication terminale. Si la deadline expire, l'abort du wrapper source détruit aussi son `JoinSet` interne et annule ses tâches imbriquées avant le terminal. Une faute déjà observée n'est pas remplacée par un stop concurrent, sauf le `drain_timeout` terminal lorsqu'une récupération bornée dépasse sa deadline. Pour Yellowstone, un timeout de fermeture Transport survenant après un Stop déjà demandé est traité comme une fermeture coopérative : `SolanaYellowstoneGrpcSubscribeSession::close()` a déjà aborté puis joint son acteur avant de retourner ce timeout. Les autres erreurs de fermeture restent terminales. L'abandon terminal d'une hydration retire son pending run-local sans le convertir artificiellement en travail `settled`.
## Admission, coalescence et backpressure

View File

@@ -1,5 +1,5 @@
<!-- file: crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md -->
<!-- version: 13 -->
<!-- version: 14 -->
# Utilisation de ksp-worker-raw-transaction-ingest-lib
@@ -46,7 +46,7 @@ Bornes publiques :
```text
admission_queue_capacity 1 ..= 65_536 défaut 256
persistence_concurrency 1 ..= 64 défaut 8
shutdown_drain_timeout 100 ms ..= 30 s défaut 5 s
shutdown_drain_timeout 100 ms ..= 30 s défaut 10 s
```
Une valeur hors borne retourne `worker_raw_transaction_ingest.settings_invalid` avec uniquement le nom stable du champ invalide.
@@ -394,7 +394,7 @@ if accepted {
}
```
`request_stop()` est idempotent. Le terminal n'est publié qu'après le drain borné et la récupération des tâches possédées. Si la deadline de drain expire, toutes les tâches source/persistence encore possédées sont abortées puis jointes avant publication terminale ; une persistence libérée après ce terminal ne peut donc pas produire une complétion tardive. Une faute source déjà observée reste prioritaire face à un stop concurrent, sauf si le drain lui-même expire et devient le terminal `drain_timeout`.
`request_stop()` est idempotent. Le terminal n'est publié qu'après le drain borné et la récupération des tâches possédées. Le défaut de drain est de 10 s afin de laisser une fermeture Yellowstone bornée à 5 s se terminer sans collision de deadline. Si la deadline de drain expire, toutes les tâches source/persistence encore possédées sont abortées puis jointes avant publication terminale ; une persistence libérée après ce terminal ne peut donc pas produire une complétion tardive. Une faute source déjà observée reste prioritaire face à un stop concurrent, sauf si le drain lui-même expire et devient le terminal `drain_timeout`.
## Interpréter les faults

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs
// version: 45
// version: 46
use sha2::Digest; // rust-rules: trait-import
@@ -401,6 +401,16 @@ impl RawTransactionIngestLiveSource {
);
}
fn kind_code(&self) -> &'static str {
return match self {
Self::HeliusTransaction(_) => "helius_transaction",
Self::HttpBlockPolling(_) => "http_block_polling",
Self::StandardBlock(_) => "standard_block",
Self::StandardLogs(_) => "standard_logs",
Self::Yellowstone(_) => "yellowstone",
};
}
fn source_key(&self) -> [u8; 32] {
return match self {
Self::HeliusTransaction(source) => source.source_key,
@@ -427,6 +437,7 @@ impl RawTransactionIngestLiveSource {
inventory_publisher: RawTransactionIngestSourceInventoryPublisher,
shared: RawTransactionIngestSourceRuntimeShared,
) -> ksp_core_lib::Result<()> {
let source_kind = self.kind_code();
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 {
@@ -453,6 +464,20 @@ impl RawTransactionIngestLiveSource {
if let std::result::Result::Err(error) = inventory_publisher.publish(source_projection_with_state(latest, terminal_state)) {
return std::result::Result::Err(error);
}
if let std::result::Result::Err(error) = &result {
let transport_domain = source_error_context_value(error, "transport_domain").unwrap_or("none");
let transport_code = source_error_context_value(error, "transport_code").unwrap_or("none");
ksp_logging_lib::warn!(
target: "ksp-worker-raw-transaction-ingest-lib",
domain = "raw_transaction_ingest.source",
source_kind = source_kind,
error_domain = error.code().domain(),
error_code = error.code().code(),
transport_domain = transport_domain,
transport_code = transport_code,
"RAW transaction ingest live source reached terminal failure"
);
}
return result;
}
changed = source_frontier_receiver.changed() => {
@@ -1224,14 +1249,14 @@ impl crate::RawTransactionIngestYellowstoneSource {
processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Failed);
return std::result::Result::Err(error);
}
return match closed {
return match resolve_yellowstone_close_after_stop(*stop_receiver.borrow(), closed) {
std::result::Result::Ok(()) => {
processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Closed);
std::result::Result::Ok(())
},
std::result::Result::Err(error) => {
processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Failed);
std::result::Result::Err(source_transport_error(error.code()))
std::result::Result::Err(error)
},
};
}
@@ -1337,14 +1362,14 @@ impl crate::RawTransactionIngestYellowstoneSource {
processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Failed);
return std::result::Result::Err(error);
}
return match closed {
return match resolve_yellowstone_close_after_stop(*stop_receiver.borrow(), closed) {
std::result::Result::Ok(()) => {
processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Closed);
std::result::Result::Ok(())
},
std::result::Result::Err(error) => {
processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Failed);
std::result::Result::Err(source_transport_error(error.code()))
std::result::Result::Err(error)
},
};
}
@@ -4338,6 +4363,23 @@ fn source_transport_error(code: ksp_core_lib::ErrorCode) -> ksp_core_lib::Error
.with_context("transport_code", code.code());
}
fn source_error_context_value<'a>(error: &'a ksp_core_lib::Error, key: &'static str) -> std::option::Option<&'a str> {
for context in error.context() {
if context.key() == key {
return std::option::Option::Some(context.value());
}
}
return std::option::Option::None;
}
fn resolve_yellowstone_close_after_stop(stop_requested: bool, closed: ksp_core_lib::Result<()>) -> ksp_core_lib::Result<()> {
return match closed {
std::result::Result::Ok(()) => std::result::Result::Ok(()),
std::result::Result::Err(error) if stop_requested && error.code() == ksp_onchain_transport_lib::ERROR_CODE_TIMEOUT => std::result::Result::Ok(()),
std::result::Result::Err(error) => std::result::Result::Err(source_transport_error(error.code())),
};
}
fn map_yellowstone_source_state(state: ksp_onchain_transport_lib::YellowstoneGrpcSubscribeState) -> crate::RawTransactionIngestSourceState {
return match state {
ksp_onchain_transport_lib::YellowstoneGrpcSubscribeState::Active => crate::RawTransactionIngestSourceState::Active,

View File

@@ -1,12 +1,12 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/settings.rs
// version: 4
// version: 5
/// Default bounded admission queue capacity for one RAW transaction ingest Worker.
pub const DEFAULT_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY: usize = 256;
/// Default number of concurrent Store persistence operations for one RAW transaction ingest Worker.
pub const DEFAULT_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY: usize = 8;
/// Default cooperative shutdown drain deadline for one RAW transaction ingest Worker.
pub const DEFAULT_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
pub const DEFAULT_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
/// Maximum bounded admission queue capacity for one RAW transaction ingest Worker.
pub const MAX_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY: usize = 65_536;
/// Maximum number of concurrent Store persistence operations for one RAW transaction ingest Worker.

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs
// version: 36
// version: 37
fn grpc_endpoint(cluster: &str) -> std::option::Option<ksp_onchain_transport_lib::YellowstoneGrpcEndpointSettings> {
return grpc_endpoint_with_identity(cluster, "yellowstone-fixture", "fixture-provider");
@@ -4614,3 +4614,25 @@ async fn v0_3_14_pre_010_repair_only_progresses_when_existing_capacity_is_one()
assert_eq!(registry.hydration_permits.available_permits(), 1);
return;
}
#[test]
fn v0_3_15_pre_009_fix_001_yellowstone_stop_accepts_only_close_timeout_as_clean_terminal() {
let timeout = ksp_core_lib::Error::new(ksp_onchain_transport_lib::ERROR_CODE_TIMEOUT, "fixture timeout");
assert!(super::resolve_yellowstone_close_after_stop(true, std::result::Result::Err(timeout)).is_ok());
let timeout_without_stop = ksp_core_lib::Error::new(ksp_onchain_transport_lib::ERROR_CODE_TIMEOUT, "fixture timeout");
let timeout_without_stop = super::resolve_yellowstone_close_after_stop(false, std::result::Result::Err(timeout_without_stop));
assert!(timeout_without_stop.is_err());
let connection = ksp_core_lib::Error::new(ksp_onchain_transport_lib::ERROR_CODE_HTTP_CONNECTION_FAILED, "fixture connection");
let connection = super::resolve_yellowstone_close_after_stop(true, std::result::Result::Err(connection));
assert!(connection.is_err());
return;
}
#[test]
fn v0_3_15_pre_009_fix_001_source_error_context_projection_is_key_bounded() {
let error = super::source_transport_error(ksp_onchain_transport_lib::ERROR_CODE_TIMEOUT);
assert_eq!(super::source_error_context_value(&error, "transport_domain"), std::option::Option::Some("onchain_transport"));
assert_eq!(super::source_error_context_value(&error, "transport_code"), std::option::Option::Some("timeout"));
assert_eq!(super::source_error_context_value(&error, "unknown"), std::option::Option::None);
return;
}

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/settings.rs
// version: 2
// version: 3
fn identities() -> std::option::Option<(ksp_store_lib::RawNetworkId, ksp_worker_api::WorkerId)> {
let network_result = ksp_store_lib::RawNetworkId::new("mainnet");
@@ -31,7 +31,7 @@ fn pre_003_defaults_are_exact_and_preserve_typed_identity() {
assert_eq!(settings.shutdown_drain_timeout(), std::time::Duration::from_secs(5));
assert_eq!(crate::DEFAULT_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY, 256);
assert_eq!(crate::DEFAULT_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY, 8);
assert_eq!(crate::DEFAULT_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT, std::time::Duration::from_secs(5));
assert_eq!(crate::DEFAULT_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT, std::time::Duration::from_secs(10));
return;
}