0.3.15-pre.009-fix.003
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs
|
||||
// version: 18
|
||||
// version: 19
|
||||
|
||||
type PersistencePort = std::sync::Arc<dyn crate::RawTransactionIngestPersistencePort + 'static>;
|
||||
type PersistenceTasks = tokio::task::JoinSet<ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome>>;
|
||||
@@ -196,7 +196,8 @@ async fn drain_admission_and_persistence(
|
||||
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
|
||||
) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
let mut fault = std::option::Option::None;
|
||||
admission.close();
|
||||
// Keep admission open while cooperative sources observe the stop signal and drop their senders.
|
||||
// Closing the receiver here races with nested source-stop propagation.
|
||||
loop {
|
||||
while persistence.len() >= settings.persistence_concurrency() {
|
||||
let joined = persistence.join_next().await;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs
|
||||
// version: 33
|
||||
// version: 34
|
||||
|
||||
//! Release-completeness canaries through the `v0.3.15-pre.004` WebSocket capability-enforcement tranche.
|
||||
|
||||
@@ -36,6 +36,7 @@ fn pre_010_production_module_inventory_is_exact() -> std::io::Result<()> {
|
||||
names,
|
||||
std::vec![
|
||||
"admission.rs",
|
||||
"constants.rs",
|
||||
"continuity.rs",
|
||||
"error.rs",
|
||||
"identity.rs",
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs
|
||||
// version: 11
|
||||
// version: 12
|
||||
|
||||
struct ActiveTaskGuard {
|
||||
active: std::sync::Arc<std::sync::atomic::AtomicUsize>,
|
||||
@@ -994,3 +994,58 @@ async fn v0_3_13_pre_011_drain_timeout_prevents_late_persistence_completion_afte
|
||||
assert_eq!(retained.in_flight_persistence(), 0);
|
||||
return;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn v0_3_15_pre_009_fix_003_shutdown_keeps_admission_open_until_cooperative_source_observes_stop() {
|
||||
let settings = match settings("mainnet") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let (started_sender, started_receiver) = tokio::sync::oneshot::channel::<()>();
|
||||
let (release_sender, release_receiver) = tokio::sync::oneshot::channel::<()>();
|
||||
let handle = match super::start_foundation_with_source_spawner(
|
||||
settings,
|
||||
tokio::runtime::Handle::current(),
|
||||
std::option::Option::None,
|
||||
move |children, mut stop_receiver, admission_sender| {
|
||||
let _abort_handle = children.spawn(async move {
|
||||
let _started = started_sender.send(());
|
||||
if release_receiver.await.is_err() {
|
||||
return std::result::Result::Err(crate::runtime_error("test.release_closed"));
|
||||
}
|
||||
if admission_sender.is_closed() {
|
||||
return std::result::Result::Err(crate::runtime_error("test.admission_closed_before_source_stop"));
|
||||
}
|
||||
if !*stop_receiver.borrow() {
|
||||
let changed = stop_receiver.changed().await;
|
||||
if changed.is_err() || !*stop_receiver.borrow() {
|
||||
return std::result::Result::Err(crate::runtime_error("test.stop_not_observed"));
|
||||
}
|
||||
}
|
||||
return std::result::Result::Ok(());
|
||||
});
|
||||
},
|
||||
) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return,
|
||||
};
|
||||
if started_receiver.await.is_err() {
|
||||
return;
|
||||
}
|
||||
assert!(handle.request_stop());
|
||||
for _ in 0..64 {
|
||||
if handle.snapshot_source().current().worker_snapshot().state() == ksp_worker_api::WorkerState::Stopping {
|
||||
break;
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
assert_eq!(handle.snapshot_source().current().worker_snapshot().state(), ksp_worker_api::WorkerState::Stopping);
|
||||
let _released = release_sender.send(());
|
||||
let terminal = match handle.wait_terminal().await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return,
|
||||
};
|
||||
assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped);
|
||||
assert_eq!(handle.snapshot_source().current().source_failure_total(), 0);
|
||||
return;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user