997 lines
45 KiB
Rust
997 lines
45 KiB
Rust
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs
|
|
// version: 11
|
|
|
|
struct ActiveTaskGuard {
|
|
active: std::sync::Arc<std::sync::atomic::AtomicUsize>,
|
|
}
|
|
|
|
impl ActiveTaskGuard {
|
|
fn new(active: std::sync::Arc<std::sync::atomic::AtomicUsize>) -> Self {
|
|
active.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
|
|
return Self { active };
|
|
}
|
|
}
|
|
|
|
impl std::ops::Drop for ActiveTaskGuard {
|
|
fn drop(&mut self) {
|
|
self.active.fetch_sub(1, std::sync::atomic::Ordering::AcqRel);
|
|
return;
|
|
}
|
|
}
|
|
|
|
fn settings(network: &str) -> std::option::Option<crate::RawTransactionIngestSettings> {
|
|
let network = match ksp_store_lib::RawNetworkId::new(network) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let worker_id = match ksp_worker_api::WorkerId::new("raw-ingest-runtime-001") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
return std::option::Option::Some(crate::RawTransactionIngestSettings::with_defaults(network, worker_id));
|
|
}
|
|
|
|
async fn wait_for_active_count(active: &std::sync::Arc<std::sync::atomic::AtomicUsize>, expected: usize) -> bool {
|
|
for _ in 0..64 {
|
|
if active.load(std::sync::atomic::Ordering::Acquire) == expected {
|
|
return true;
|
|
}
|
|
tokio::task::yield_now().await;
|
|
}
|
|
return active.load(std::sync::atomic::Ordering::Acquire) == expected;
|
|
}
|
|
|
|
#[test]
|
|
fn pre_004_current_runtime_is_required_before_spawn() {
|
|
let result = super::current_runtime_handle();
|
|
assert!(result.is_err());
|
|
let error = match result {
|
|
std::result::Result::Ok(_) => return,
|
|
std::result::Result::Err(value) => value,
|
|
};
|
|
assert_eq!(error.code(), crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
|
assert_eq!(error.context().len(), 1);
|
|
assert_eq!(error.context()[0].value(), "start.runtime_unavailable");
|
|
return;
|
|
}
|
|
|
|
#[test]
|
|
fn pre_004_store_network_mismatch_is_rejected_without_echoing_network_values() {
|
|
let settings = match settings("mainnet") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let other = match ksp_store_lib::RawNetworkId::new("devnet") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
let result = super::validate_store_network(&settings, &other);
|
|
assert!(result.is_err());
|
|
let error = match result {
|
|
std::result::Result::Ok(_) => return,
|
|
std::result::Result::Err(value) => value,
|
|
};
|
|
assert_eq!(error.code(), crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
|
let rendered = std::format!("{error:?}");
|
|
assert!(!rendered.contains("mainnet"));
|
|
assert!(!rendered.contains("devnet"));
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_004_immediate_stop_is_idempotent_and_reaches_stopped_terminal() {
|
|
let settings = match settings("mainnet") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let handle = match super::start_foundation(settings, tokio::runtime::Handle::current(), std::option::Option::None) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(handle.snapshot_source().current().worker_snapshot().state(), ksp_worker_api::WorkerState::Starting);
|
|
assert!(handle.request_stop());
|
|
assert!(!handle.request_stop());
|
|
let terminal = handle.wait_terminal().await;
|
|
assert!(terminal.is_ok());
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped);
|
|
assert!(terminal.is_terminal());
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_004_running_lifecycle_stops_through_cloned_handle() {
|
|
let settings = match settings("mainnet") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let handle = match super::start_foundation(settings, tokio::runtime::Handle::current(), std::option::Option::None) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
tokio::task::yield_now().await;
|
|
assert_eq!(handle.snapshot_source().current().worker_snapshot().state(), ksp_worker_api::WorkerState::Running);
|
|
let clone = handle.clone();
|
|
assert!(clone.request_stop());
|
|
assert!(!handle.request_stop());
|
|
let terminal = handle.wait_terminal().await;
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_004_dropping_last_control_handle_causes_private_runtime_exit() {
|
|
let settings = match settings("mainnet") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let handle = match super::start_foundation(settings, tokio::runtime::Handle::current(), std::option::Option::None) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
let observer = handle.snapshot_source();
|
|
std::mem::drop(handle);
|
|
loop {
|
|
let current = observer.current();
|
|
if current.worker_snapshot().state().is_terminal() {
|
|
assert_eq!(current.worker_snapshot().state(), ksp_worker_api::WorkerState::Stopped);
|
|
return;
|
|
}
|
|
let _changed = observer.wait_for_change(current.worker_snapshot().sequence()).await;
|
|
}
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_005_supervisor_joins_all_cooperative_source_tasks_before_terminal() {
|
|
let settings = match settings("mainnet") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
|
let source_active = std::sync::Arc::clone(&active);
|
|
let handle = match super::start_foundation_with_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::None,
|
|
move |children, stop_receiver, _admission_sender| {
|
|
for _ in 0..3 {
|
|
let active = std::sync::Arc::clone(&source_active);
|
|
let mut child_stop = stop_receiver.clone();
|
|
let _abort_handle = children.spawn(async move {
|
|
let _guard = ActiveTaskGuard::new(active);
|
|
loop {
|
|
let changed = child_stop.changed().await;
|
|
if changed.is_err() || *child_stop.borrow() {
|
|
return std::result::Result::Ok(());
|
|
}
|
|
}
|
|
});
|
|
}
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert!(wait_for_active_count(&active, 3).await);
|
|
assert!(handle.request_stop());
|
|
let terminal = handle.wait_terminal().await;
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped);
|
|
assert_eq!(active.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_005_supervisor_reaps_completed_source_task_and_still_joins_live_child_on_stop() {
|
|
let settings = match settings("mainnet") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
|
let completed = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
|
let source_active = std::sync::Arc::clone(&active);
|
|
let source_completed = std::sync::Arc::clone(&completed);
|
|
let handle = match super::start_foundation_with_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::None,
|
|
move |children, stop_receiver, _admission_sender| {
|
|
let completed = std::sync::Arc::clone(&source_completed);
|
|
let _completed_abort_handle = children.spawn(async move {
|
|
completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
|
|
return std::result::Result::Ok(());
|
|
});
|
|
let active = std::sync::Arc::clone(&source_active);
|
|
let mut child_stop = stop_receiver.clone();
|
|
let _active_abort_handle = children.spawn(async move {
|
|
let _guard = ActiveTaskGuard::new(active);
|
|
loop {
|
|
let changed = child_stop.changed().await;
|
|
if changed.is_err() || *child_stop.borrow() {
|
|
return std::result::Result::Ok(());
|
|
}
|
|
}
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert!(wait_for_active_count(&active, 1).await);
|
|
for _ in 0..64 {
|
|
if completed.load(std::sync::atomic::Ordering::Acquire) == 1 {
|
|
break;
|
|
}
|
|
tokio::task::yield_now().await;
|
|
}
|
|
assert_eq!(completed.load(std::sync::atomic::Ordering::Acquire), 1);
|
|
assert!(handle.request_stop());
|
|
let terminal = handle.wait_terminal().await;
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped);
|
|
assert_eq!(active.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
return;
|
|
}
|
|
|
|
#[derive(Clone, Copy)]
|
|
enum RuntimePortResponse {
|
|
Conflict,
|
|
StoreFailure,
|
|
StoreFailureBlocked,
|
|
SuccessBlocked,
|
|
}
|
|
|
|
struct RuntimePortActiveGuard<'a> {
|
|
active: &'a std::sync::atomic::AtomicUsize,
|
|
}
|
|
|
|
impl std::ops::Drop for RuntimePortActiveGuard<'_> {
|
|
fn drop(&mut self) {
|
|
self.active.fetch_sub(1, std::sync::atomic::Ordering::AcqRel);
|
|
return;
|
|
}
|
|
}
|
|
|
|
struct RuntimePersistencePort {
|
|
active: std::sync::atomic::AtomicUsize,
|
|
completed: std::sync::atomic::AtomicUsize,
|
|
max_active: std::sync::atomic::AtomicUsize,
|
|
network: ksp_store_lib::RawNetworkId,
|
|
normal_mode_seen: std::sync::atomic::AtomicBool,
|
|
notify: tokio::sync::Notify,
|
|
released: std::sync::atomic::AtomicBool,
|
|
response: RuntimePortResponse,
|
|
}
|
|
|
|
impl RuntimePersistencePort {
|
|
fn new(network: ksp_store_lib::RawNetworkId, response: RuntimePortResponse, released: bool) -> Self {
|
|
return Self {
|
|
active: std::sync::atomic::AtomicUsize::new(0),
|
|
completed: std::sync::atomic::AtomicUsize::new(0),
|
|
max_active: std::sync::atomic::AtomicUsize::new(0),
|
|
network,
|
|
normal_mode_seen: std::sync::atomic::AtomicBool::new(false),
|
|
notify: tokio::sync::Notify::new(),
|
|
released: std::sync::atomic::AtomicBool::new(released),
|
|
response,
|
|
};
|
|
}
|
|
|
|
fn release(&self) {
|
|
self.released.store(true, std::sync::atomic::Ordering::Release);
|
|
self.notify.notify_waiters();
|
|
return;
|
|
}
|
|
|
|
fn update_max_active(&self, current: usize) {
|
|
let mut observed = self.max_active.load(std::sync::atomic::Ordering::Acquire);
|
|
while current > observed {
|
|
match self.max_active.compare_exchange_weak(observed, current, std::sync::atomic::Ordering::AcqRel, std::sync::atomic::Ordering::Acquire) {
|
|
std::result::Result::Ok(_) => return,
|
|
std::result::Result::Err(value) => observed = value,
|
|
}
|
|
}
|
|
return;
|
|
}
|
|
}
|
|
|
|
impl crate::RawTransactionIngestPersistencePort for RuntimePersistencePort {
|
|
fn network_matches(&self, network: &ksp_store_lib::RawNetworkId) -> bool {
|
|
return &self.network == network;
|
|
}
|
|
|
|
fn persist_acquisition<'a>(
|
|
&'a self,
|
|
_transaction: ksp_store_lib::RawTransaction,
|
|
_observation: ksp_store_lib::RawTransactionObservation,
|
|
mode: ksp_store_lib::RawTransactionAcquisitionMode,
|
|
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawAcquisitionWriteOutcome>> {
|
|
self.normal_mode_seen.store(mode == ksp_store_lib::RawTransactionAcquisitionMode::Normal, std::sync::atomic::Ordering::Release);
|
|
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::RawAcquisitionWriteOutcome::new(
|
|
ksp_store_lib::RawEntityWriteOutcome::Inserted,
|
|
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 record_observation<'a>(
|
|
&'a self,
|
|
_observation: ksp_store_lib::RawTransactionObservation,
|
|
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawObservationWriteOutcome>> {
|
|
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<crate::RawTransactionIngress> {
|
|
let signature = ksp_store_lib::RawTransactionSignature::new([signature_byte; 64]);
|
|
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64(
|
|
network.clone(),
|
|
signature,
|
|
42,
|
|
std::option::Option::Some(1_700_000_000),
|
|
"AQID",
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
);
|
|
let provider = match ksp_store_lib::RawProvenanceCode::new("deterministic-harness") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let protocol = match ksp_store_lib::RawProvenanceCode::new("internal") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let method = match ksp_store_lib::RawProvenanceCode::new("runtime") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let received_at = match ksp_store_lib::RawTimestamp::from_unix_millis(1_700_000_000_000) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let provenance = ksp_store_lib::RawAcquisitionProvenance::new(provider, protocol, method, ksp_store_lib::RawAcquisitionOrigin::Live, received_at);
|
|
return std::option::Option::Some(crate::RawTransactionIngress { material, network: network.clone(), provenance, source_key: [signature_byte; 32] });
|
|
}
|
|
|
|
fn settings_with_persistence_concurrency(concurrency: usize) -> std::option::Option<crate::RawTransactionIngestSettings> {
|
|
return settings_with_runtime_limits(8, concurrency, std::time::Duration::from_secs(5));
|
|
}
|
|
|
|
fn settings_with_runtime_limits(
|
|
admission_queue_capacity: usize,
|
|
persistence_concurrency: usize,
|
|
shutdown_drain_timeout: std::time::Duration,
|
|
) -> std::option::Option<crate::RawTransactionIngestSettings> {
|
|
let network = match ksp_store_lib::RawNetworkId::new("mainnet") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let worker_id = match ksp_worker_api::WorkerId::new("raw-ingest-runtime-persistence-001") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
return match crate::RawTransactionIngestSettings::new(network, worker_id, admission_queue_capacity, persistence_concurrency, shutdown_drain_timeout) {
|
|
std::result::Result::Ok(value) => std::option::Option::Some(value),
|
|
std::result::Result::Err(_) => std::option::Option::None,
|
|
};
|
|
}
|
|
|
|
async fn wait_for_completed(port: &std::sync::Arc<RuntimePersistencePort>, expected: usize) -> bool {
|
|
for _ in 0..128 {
|
|
if port.completed.load(std::sync::atomic::Ordering::Acquire) == expected {
|
|
return true;
|
|
}
|
|
tokio::task::yield_now().await;
|
|
}
|
|
return port.completed.load(std::sync::atomic::Ordering::Acquire) == expected;
|
|
}
|
|
|
|
async fn wait_for_max_active(port: &std::sync::Arc<RuntimePersistencePort>, expected: usize) -> bool {
|
|
for _ in 0..128 {
|
|
if port.max_active.load(std::sync::atomic::Ordering::Acquire) == expected {
|
|
return true;
|
|
}
|
|
tokio::task::yield_now().await;
|
|
}
|
|
return port.max_active.load(std::sync::atomic::Ordering::Acquire) == expected;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_007_runtime_bounds_in_flight_store_persistence_to_configured_concurrency() {
|
|
let settings = match settings_with_persistence_concurrency(2) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, false));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
for signature_byte in 1..=4 {
|
|
let ingress = match runtime_ingress(&network, signature_byte) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
}
|
|
return std::result::Result::Ok(());
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert!(wait_for_max_active(&port, 2).await);
|
|
assert_eq!(port.active.load(std::sync::atomic::Ordering::Acquire), 2);
|
|
assert_eq!(port.max_active.load(std::sync::atomic::Ordering::Acquire), 2);
|
|
port.release();
|
|
assert!(wait_for_completed(&port, 4).await);
|
|
assert!(port.normal_mode_seen.load(std::sync::atomic::Ordering::Acquire));
|
|
assert_eq!(port.max_active.load(std::sync::atomic::Ordering::Acquire), 2);
|
|
assert!(handle.request_stop());
|
|
let terminal = handle.wait_terminal().await;
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped);
|
|
assert_eq!(port.active.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_007_content_conflict_becomes_terminal_after_private_drain() {
|
|
let settings = match settings_with_persistence_concurrency(1) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::Conflict, true));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
let ingress = match runtime_ingress(&network, 11) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
return std::result::Result::Ok(());
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
let source = handle.snapshot_source();
|
|
let terminal = handle.wait_terminal().await;
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT));
|
|
assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 1);
|
|
let snapshot = source.current();
|
|
assert_eq!(snapshot.content_conflict_total(), 1);
|
|
assert_eq!(snapshot.store_failure_total(), 0);
|
|
assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_007_store_failure_becomes_terminal_without_exposing_remote_error_text() {
|
|
let settings = match settings_with_persistence_concurrency(1) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::StoreFailure, true));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
let ingress = match runtime_ingress(&network, 12) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
return std::result::Result::Ok(());
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
let terminal = handle.wait_terminal().await;
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED));
|
|
assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 1);
|
|
assert!(!std::format!("{handle:?}").contains("synthetic store failure"));
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_008_runtime_snapshots_count_successful_pipeline_and_retain_terminal() {
|
|
let settings = match settings_with_persistence_concurrency(2) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, true));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
for signature_byte in 21..=23 {
|
|
let ingress = match runtime_ingress(&network, signature_byte) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
}
|
|
return std::result::Result::Ok(());
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
let concrete_source = handle.snapshot_source();
|
|
let common_source = handle.worker_snapshot_source();
|
|
assert!(wait_for_completed(&port, 3).await);
|
|
assert!(handle.request_stop());
|
|
let terminal = handle.wait_terminal().await;
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped);
|
|
let snapshot = concrete_source.current();
|
|
assert_eq!(snapshot.worker_snapshot().state(), ksp_worker_api::WorkerState::Stopped);
|
|
assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Healthy);
|
|
assert_eq!(snapshot.worker_snapshot().activity(), ksp_worker_api::WorkerActivity::Idle);
|
|
assert_eq!(snapshot.admitted_total(), 3);
|
|
assert_eq!(snapshot.canonicalized_total(), 3);
|
|
assert_eq!(snapshot.persisted_total(), 3);
|
|
assert_eq!(snapshot.entity_inserted_total(), 3);
|
|
assert_eq!(snapshot.entity_already_present_total(), 0);
|
|
assert_eq!(snapshot.entity_skipped_purged_total(), 0);
|
|
assert_eq!(snapshot.observation_inserted_total(), 3);
|
|
assert_eq!(snapshot.observation_already_present_total(), 0);
|
|
assert_eq!(snapshot.content_conflict_total(), 0);
|
|
assert_eq!(snapshot.store_failure_total(), 0);
|
|
assert_eq!(snapshot.source_failure_total(), 0);
|
|
assert_eq!(snapshot.backpressure_wait_total(), 0);
|
|
assert_eq!(snapshot.admission_queue_depth(), 0);
|
|
assert_eq!(snapshot.in_flight_persistence(), 0);
|
|
let common = ksp_worker_api::WorkerSnapshotSource::current(&common_source);
|
|
assert_eq!(&common, snapshot.worker_snapshot());
|
|
let retained = concrete_source.wait_for_change(snapshot.worker_snapshot().sequence()).await;
|
|
assert_eq!(retained, snapshot);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_008_fault_snapshots_count_classified_store_failures_and_project_unhealthy() {
|
|
let settings = match settings_with_persistence_concurrency(1) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::StoreFailure, true));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
let ingress = match runtime_ingress(&network, 24) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
return std::result::Result::Ok(());
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
let source = handle.snapshot_source();
|
|
let terminal = handle.wait_terminal().await;
|
|
let terminal = match terminal {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(terminal, ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED));
|
|
let snapshot = source.current();
|
|
assert_eq!(snapshot.worker_snapshot().state(), terminal);
|
|
assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy);
|
|
assert_eq!(snapshot.worker_snapshot().activity(), ksp_worker_api::WorkerActivity::Idle);
|
|
assert_eq!(snapshot.admitted_total(), 1);
|
|
assert_eq!(snapshot.canonicalized_total(), 1);
|
|
assert_eq!(snapshot.persisted_total(), 0);
|
|
assert_eq!(snapshot.store_failure_total(), 1);
|
|
assert_eq!(snapshot.content_conflict_total(), 0);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_009_source_failure_is_counted_and_late_stop_cannot_replace_terminal_fault() {
|
|
let settings = match settings("mainnet") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let handle = match super::start_foundation_with_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::None,
|
|
move |children: &mut tokio::task::JoinSet<ksp_core_lib::Result<()>>, _stop_receiver, _admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
let source = handle.snapshot_source();
|
|
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::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_SOURCE_FAILED));
|
|
let snapshot = source.current();
|
|
assert_eq!(snapshot.source_failure_total(), 1);
|
|
assert_eq!(snapshot.backpressure_wait_total(), 0);
|
|
assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy);
|
|
assert!(!handle.request_stop());
|
|
let retained = match handle.wait_terminal().await {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(retained, terminal);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_009_stop_does_not_hide_store_failure_observed_during_bounded_drain() {
|
|
let settings = match settings_with_persistence_concurrency(1) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::StoreFailureBlocked, false));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
let ingress = match runtime_ingress(&network, 41) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.ingress_invalid")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
return std::result::Result::Ok(());
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert!(wait_for_max_active(&port, 1).await);
|
|
assert!(handle.request_stop());
|
|
port.release();
|
|
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::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED));
|
|
let snapshot = handle.snapshot_source().current();
|
|
assert_eq!(snapshot.store_failure_total(), 1);
|
|
assert_eq!(snapshot.in_flight_persistence(), 0);
|
|
assert_eq!(port.active.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_009_saturation_is_observable_without_drop_or_unbounded_admission() {
|
|
let settings = match settings_with_runtime_limits(1, 1, std::time::Duration::from_secs(5)) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, false));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let source_stage = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
|
let source_stage_for_task = source_stage.clone();
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
for signature_byte in 51..=53 {
|
|
let ingress = match runtime_ingress(&network, signature_byte) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.ingress_invalid")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
source_stage_for_task.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
|
|
}
|
|
return std::result::Result::Ok(());
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert!(wait_for_max_active(&port, 1).await);
|
|
for _ in 0..128 {
|
|
if source_stage.load(std::sync::atomic::Ordering::Acquire) == 2 {
|
|
break;
|
|
}
|
|
tokio::task::yield_now().await;
|
|
}
|
|
assert_eq!(source_stage.load(std::sync::atomic::Ordering::Acquire), 2);
|
|
port.release();
|
|
assert!(wait_for_completed(&port, 3).await);
|
|
for _ in 0..128 {
|
|
if source_stage.load(std::sync::atomic::Ordering::Acquire) == 3 {
|
|
break;
|
|
}
|
|
tokio::task::yield_now().await;
|
|
}
|
|
assert_eq!(source_stage.load(std::sync::atomic::Ordering::Acquire), 3);
|
|
assert!(handle.request_stop());
|
|
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);
|
|
let snapshot = handle.snapshot_source().current();
|
|
assert_eq!(snapshot.persisted_total(), 3);
|
|
assert_eq!(snapshot.source_failure_total(), 0);
|
|
assert!(snapshot.backpressure_wait_total() >= 1);
|
|
assert_eq!(snapshot.admission_queue_depth(), 0);
|
|
assert_eq!(snapshot.in_flight_persistence(), 0);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_009_drain_timeout_aborts_and_joins_all_owned_source_and_persistence_tasks() {
|
|
let settings = match settings_with_runtime_limits(1, 1, crate::MIN_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, false));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let source_active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
|
let source_active_for_task = source_active.clone();
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
let ingress = match runtime_ingress(&network, 61) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.ingress_invalid")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
let _guard = ActiveTaskGuard::new(source_active_for_task);
|
|
return std::future::pending::<ksp_core_lib::Result<()>>().await;
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert!(wait_for_max_active(&port, 1).await);
|
|
assert!(wait_for_active_count(&source_active, 1).await);
|
|
assert!(handle.request_stop());
|
|
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::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_DRAIN_TIMEOUT));
|
|
let snapshot = handle.snapshot_source().current();
|
|
assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy);
|
|
assert_eq!(snapshot.admission_queue_depth(), 0);
|
|
assert_eq!(snapshot.in_flight_persistence(), 0);
|
|
assert_eq!(source_active.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
assert_eq!(port.active.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn v0_3_13_pre_011_source_counter_exhaustion_stays_terminal_and_is_not_collapsed_to_source_failed() {
|
|
let settings = match settings("mainnet") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let handle = match super::start_foundation_with_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::None,
|
|
move |children: &mut tokio::task::JoinSet<ksp_core_lib::Result<()>>, _stop_receiver, _admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
return std::result::Result::Err(crate::counter_exhausted_error("source_reconnect_total"));
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
let source = handle.snapshot_source();
|
|
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::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_COUNTER_EXHAUSTED));
|
|
let snapshot = source.current();
|
|
assert_eq!(snapshot.source_failure_total(), 1);
|
|
assert_eq!(snapshot.worker_snapshot().state(), terminal);
|
|
assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy);
|
|
let terminal_sequence = snapshot.worker_snapshot().sequence();
|
|
assert!(!handle.request_stop());
|
|
let retained = source.current();
|
|
assert_eq!(retained.worker_snapshot().state(), terminal);
|
|
assert_eq!(retained.worker_snapshot().sequence(), terminal_sequence);
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn v0_3_13_pre_011_drain_timeout_prevents_late_persistence_completion_after_terminal() {
|
|
let settings = match settings_with_runtime_limits(1, 1, crate::MIN_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let network = settings.network().clone();
|
|
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, false));
|
|
let runtime_port: super::PersistencePort = port.clone();
|
|
let source_active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
|
let source_active_for_task = std::sync::Arc::clone(&source_active);
|
|
let handle = match super::start_foundation_with_port_and_source_spawner(
|
|
settings,
|
|
tokio::runtime::Handle::current(),
|
|
std::option::Option::Some(runtime_port),
|
|
move |children, _stop_receiver, admission_sender| {
|
|
let _abort_handle = children.spawn(async move {
|
|
let ingress = match runtime_ingress(&network, 81) {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.ingress_invalid")),
|
|
};
|
|
if admission_sender.send(ingress).await.is_err() {
|
|
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
|
|
}
|
|
let _guard = ActiveTaskGuard::new(source_active_for_task);
|
|
return std::future::pending::<ksp_core_lib::Result<()>>().await;
|
|
});
|
|
},
|
|
) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert!(wait_for_max_active(&port, 1).await);
|
|
assert!(wait_for_active_count(&source_active, 1).await);
|
|
assert!(handle.request_stop());
|
|
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::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_DRAIN_TIMEOUT));
|
|
let snapshot = handle.snapshot_source().current();
|
|
let terminal_sequence = snapshot.worker_snapshot().sequence();
|
|
assert_eq!(snapshot.in_flight_persistence(), 0);
|
|
assert_eq!(port.active.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
assert_eq!(source_active.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
port.release();
|
|
for _ in 0..64 {
|
|
tokio::task::yield_now().await;
|
|
}
|
|
assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 0);
|
|
let retained = handle.snapshot_source().current();
|
|
assert_eq!(retained.worker_snapshot().state(), terminal);
|
|
assert_eq!(retained.worker_snapshot().sequence(), terminal_sequence);
|
|
assert_eq!(retained.in_flight_persistence(), 0);
|
|
return;
|
|
}
|