v0.3.14-pre.012

This commit is contained in:
2026-09-12 20:28:00 +02:00
parent f3c4169752
commit 8a0717c839
7 changed files with 499 additions and 27 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs
// version: 42
// version: 43
use sha2::Digest; // rust-rules: trait-import
@@ -1726,6 +1726,9 @@ impl crate::RawTransactionIngestHttpBlockPollingSource {
break 'source;
}
}
if *stop_receiver.borrow() {
break;
}
if !blocked_by_null && let std::option::Option::Some(proven_end_slot) = discovery.proven_end_slot {
let coverage_result = {
let mut contracts = match continuity_contracts.lock() {
@@ -2584,6 +2587,7 @@ impl crate::RawTransactionIngestRuntimeResources {
std::sync::Arc::clone(&continuity_contracts),
std::sync::Arc::clone(&inventory),
processing_frontier_sender,
settings.shutdown_drain_timeout(),
)
.await;
let validation = {
@@ -2682,18 +2686,19 @@ async fn supervise_live_source_tasks(
continuity_contracts: std::sync::Arc<std::sync::Mutex<crate::RawTransactionIngestContinuityContracts>>,
inventory: std::sync::Arc<std::sync::Mutex<RawTransactionIngestSourceInventory>>,
processing_frontier_sender: tokio::sync::watch::Sender<crate::RawTransactionIngestProcessingFrontierProjection>,
shutdown_drain_timeout: std::time::Duration,
) -> ksp_core_lib::Result<()> {
loop {
if *stop_receiver.borrow() {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::None).await;
return drain_live_source_tasks(&mut children, std::option::Option::None, shutdown_drain_timeout).await;
}
let joined = tokio::select! {
biased;
changed = stop_receiver.changed() => {
if changed.is_err() || *stop_receiver.borrow() {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::None).await;
return drain_live_source_tasks(&mut children, std::option::Option::None, shutdown_drain_timeout).await;
}
continue;
}
@@ -2703,11 +2708,21 @@ async fn supervise_live_source_tasks(
std::option::Option::Some(std::result::Result::Ok(value)) => value,
std::option::Option::Some(std::result::Result::Err(_)) => {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(crate::runtime_error("source.task_join_failed"))).await;
return drain_live_source_tasks(
&mut children,
std::option::Option::Some(crate::runtime_error("source.task_join_failed")),
shutdown_drain_timeout,
)
.await;
},
std::option::Option::None => {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(crate::runtime_error("source.task_set_empty"))).await;
return drain_live_source_tasks(
&mut children,
std::option::Option::Some(crate::runtime_error("source.task_set_empty")),
shutdown_drain_timeout,
)
.await;
},
};
let source_fault = match source_result {
@@ -2716,7 +2731,7 @@ async fn supervise_live_source_tasks(
};
if !source_loss_is_reconcilable(&source_fault) {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(source_fault)).await;
return drain_live_source_tasks(&mut children, std::option::Option::Some(source_fault), shutdown_drain_timeout).await;
}
let supervisor_state = {
let inventory = match inventory.lock() {
@@ -2729,21 +2744,21 @@ async fn supervise_live_source_tasks(
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(error)).await;
return drain_live_source_tasks(&mut children, std::option::Option::Some(error), shutdown_drain_timeout).await;
},
};
let continuity_range = match source_loss_continuity_range(&source_fault) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(error)).await;
return drain_live_source_tasks(&mut children, std::option::Option::Some(error), shutdown_drain_timeout).await;
},
};
let continuity_range = match continuity_range {
std::option::Option::Some(value) => value,
std::option::Option::None => {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(source_fault)).await;
return drain_live_source_tasks(&mut children, std::option::Option::Some(source_fault), shutdown_drain_timeout).await;
},
};
let decision = {
@@ -2761,7 +2776,7 @@ async fn supervise_live_source_tasks(
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(error)).await;
return drain_live_source_tasks(&mut children, std::option::Option::Some(error), shutdown_drain_timeout).await;
},
};
match decision {
@@ -2770,7 +2785,7 @@ async fn supervise_live_source_tasks(
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(error)).await;
return drain_live_source_tasks(&mut children, std::option::Option::Some(error), shutdown_drain_timeout).await;
},
};
processing_frontier_sender.send_replace(aggregate);
@@ -2778,7 +2793,7 @@ async fn supervise_live_source_tasks(
},
crate::RawTransactionIngestSourceLossDecision::Fault => {
source_stop_sender.send_replace(true);
return drain_live_source_tasks(&mut children, std::option::Option::Some(source_fault)).await;
return drain_live_source_tasks(&mut children, std::option::Option::Some(source_fault), shutdown_drain_timeout).await;
},
}
}
@@ -2834,16 +2849,29 @@ fn source_loss_is_reconcilable(error: &ksp_core_lib::Error) -> bool {
async fn drain_live_source_tasks(
children: &mut tokio::task::JoinSet<([u8; 32], ksp_core_lib::Result<()>)>,
mut first_fault: std::option::Option<ksp_core_lib::Error>,
shutdown_drain_timeout: std::time::Duration,
) -> ksp_core_lib::Result<()> {
while let std::option::Option::Some(joined) = children.join_next().await {
if first_fault.is_some() {
continue;
let drain = async {
while let std::option::Option::Some(joined) = children.join_next().await {
if first_fault.is_some() {
continue;
}
first_fault = match joined {
std::result::Result::Ok((_source_key, std::result::Result::Ok(()))) => std::option::Option::None,
std::result::Result::Ok((_source_key, std::result::Result::Err(error))) => std::option::Option::Some(error),
std::result::Result::Err(_) => std::option::Option::Some(crate::runtime_error("source.task_join_failed")),
};
}
};
if tokio::time::timeout(shutdown_drain_timeout, drain).await.is_err() {
children.abort_all();
while children.join_next().await.is_some() {}
if first_fault.is_none() {
first_fault = std::option::Option::Some(ksp_core_lib::Error::new(
crate::ERROR_CODE_RAW_TRANSACTION_INGEST_DRAIN_TIMEOUT,
"RAW transaction ingest source drain timed out",
));
}
first_fault = match joined {
std::result::Result::Ok((_source_key, std::result::Result::Ok(()))) => std::option::Option::None,
std::result::Result::Ok((_source_key, std::result::Result::Err(error))) => std::option::Option::Some(error),
std::result::Result::Err(_) => std::option::Option::Some(crate::runtime_error("source.task_join_failed")),
};
}
return match first_fault {
std::option::Option::Some(error) => std::result::Result::Err(error),
@@ -4656,6 +4684,9 @@ impl RawTransactionIngestHydrationCoordinator {
std::result::Result::Ok(std::result::Result::Err(error)) => return std::result::Result::Err(error),
std::result::Result::Err(_) => return std::result::Result::Err(crate::runtime_error("source.hydration_task_join_failed")),
};
if *stop_receiver.borrow() {
return std::result::Result::Ok(false);
}
let pending = match self.pending.remove(&fetched.key) {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("source.hydration_result_without_pending")),
@@ -4665,6 +4696,9 @@ impl RawTransactionIngestHydrationCoordinator {
}
self.pending_signal_count -= pending.signals.len();
for pending_signal in pending.signals {
if *stop_receiver.borrow() {
return std::result::Result::Ok(false);
}
let signal_slot = pending_signal.signal.slot;
let resolution = resolve_known_reference_hydration(hydration, settings, pending_signal.signal, pending_signal.received_at, &fetched.observed);
let resolution = match resolution {