v0.3.11-pre.007
This commit is contained in:
@@ -1,12 +1,23 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/error.rs
|
||||
// version: 3
|
||||
// version: 4
|
||||
|
||||
/// Error code used when durable Store content conflicts with one admitted canonical RAW transaction.
|
||||
pub const ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT: ksp_core_lib::ErrorCode =
|
||||
ksp_core_lib::ErrorCode::new("worker_raw_transaction_ingest", "content_conflict");
|
||||
/// Error code used when the RAW transaction ingest Worker reaches an invalid runtime or lifecycle condition.
|
||||
pub const ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID: ksp_core_lib::ErrorCode =
|
||||
ksp_core_lib::ErrorCode::new("worker_raw_transaction_ingest", "runtime_invalid");
|
||||
/// Error code used when RAW transaction ingest Worker settings violate one bounded runtime invariant.
|
||||
pub const ERROR_CODE_RAW_TRANSACTION_INGEST_SETTINGS_INVALID: ksp_core_lib::ErrorCode =
|
||||
ksp_core_lib::ErrorCode::new("worker_raw_transaction_ingest", "settings_invalid");
|
||||
/// Error code used when Store persistence fails for one admitted canonical RAW transaction.
|
||||
pub const ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED: ksp_core_lib::ErrorCode =
|
||||
ksp_core_lib::ErrorCode::new("worker_raw_transaction_ingest", "store_failed");
|
||||
|
||||
/// Creates one terminal content-conflict error without copying conflicting material into diagnostics.
|
||||
pub(crate) fn content_conflict_error() -> ksp_core_lib::Error {
|
||||
return ksp_core_lib::Error::new(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT, "RAW transaction ingest Store content conflict");
|
||||
}
|
||||
|
||||
/// Creates one runtime-domain error carrying only one stable internal condition code.
|
||||
pub(crate) fn runtime_error(condition: &'static str) -> ksp_core_lib::Error {
|
||||
@@ -19,3 +30,10 @@ pub(crate) fn settings_error(field: &'static str) -> ksp_core_lib::Error {
|
||||
return ksp_core_lib::Error::new(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_SETTINGS_INVALID, "invalid RAW transaction ingest Worker settings")
|
||||
.with_context("field", field);
|
||||
}
|
||||
|
||||
/// Creates one terminal Store error while retaining only the already-safe lower-layer ErrorCode.
|
||||
pub(crate) fn store_error(store_code: ksp_core_lib::ErrorCode) -> ksp_core_lib::Error {
|
||||
return ksp_core_lib::Error::new(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED, "RAW transaction ingest Store persistence failed")
|
||||
.with_context("store_domain", store_code.domain())
|
||||
.with_context("store_code", store_code.code());
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs
|
||||
// version: 5
|
||||
// version: 6
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -9,19 +9,25 @@
|
||||
//!
|
||||
//! This tranche owns the concrete Worker family identity, validated technical settings
|
||||
//! and the caller-runtime-owned lifecycle with private child-task supervision. This tranche also
|
||||
//! owns bounded source-neutral admission plus common RAW canonicalization/assembly; persistence and
|
||||
//! latest-value snapshots remain in later prereleases, and no live source or Transport dependency exists.
|
||||
//! owns bounded source-neutral admission, common RAW canonicalization/assembly and backend-neutral
|
||||
//! Store persistence in normal mode; latest-value snapshots remain in a later prerelease, and no live
|
||||
//! source or Transport dependency exists.
|
||||
|
||||
mod admission;
|
||||
mod error;
|
||||
mod identity;
|
||||
mod persistence;
|
||||
mod runtime;
|
||||
mod settings;
|
||||
|
||||
/// Error code used when durable Store content conflicts with one admitted canonical RAW transaction.
|
||||
pub use self::error::ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT;
|
||||
/// Error code used when the RAW transaction ingest Worker reaches an invalid runtime or lifecycle condition.
|
||||
pub use self::error::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID;
|
||||
/// Error code used when RAW transaction ingest Worker settings violate one bounded runtime invariant.
|
||||
pub use self::error::ERROR_CODE_RAW_TRANSACTION_INGEST_SETTINGS_INVALID;
|
||||
/// Error code used when Store persistence fails for one admitted canonical RAW transaction.
|
||||
pub use self::error::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED;
|
||||
/// Stable Worker kind code used by the continuous RAW transaction ingest vertical.
|
||||
pub use self::identity::RAW_TRANSACTION_INGEST_WORKER_KIND_CODE;
|
||||
/// Cloneable external control handle for one continuous RAW transaction ingest Worker.
|
||||
@@ -55,7 +61,21 @@ pub use self::settings::RawTransactionIngestSettings;
|
||||
pub(crate) use self::admission::RawTransactionAdmission;
|
||||
/// Crate-private source-neutral ingress sent through the bounded central admission queue.
|
||||
pub(crate) use self::admission::RawTransactionIngress;
|
||||
/// Creates one terminal content-conflict error without copying conflicting material into diagnostics.
|
||||
pub(crate) use self::error::content_conflict_error;
|
||||
/// Creates one runtime-domain error without copying runtime/provider/Store values into diagnostics.
|
||||
pub(crate) use self::error::runtime_error;
|
||||
/// Creates one settings-domain error without copying caller-supplied values into diagnostics.
|
||||
pub(crate) use self::error::settings_error;
|
||||
/// Creates one terminal Store error while retaining only the already-safe lower-layer ErrorCode.
|
||||
pub(crate) use self::error::store_error;
|
||||
/// Canonical entity disposition produced by one successful Worker Store persistence attempt.
|
||||
pub(crate) use self::persistence::RawTransactionIngestEntityPersistence;
|
||||
/// Observation disposition produced by one successful Worker Store persistence attempt.
|
||||
pub(crate) use self::persistence::RawTransactionIngestObservationPersistence;
|
||||
/// 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;
|
||||
|
||||
138
crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs
Normal file
138
crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs
Normal file
@@ -0,0 +1,138 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs
|
||||
// version: 1
|
||||
|
||||
/// Canonical entity disposition produced by one Worker Store persistence attempt.
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub(crate) enum RawTransactionIngestEntityPersistence {
|
||||
/// The canonical RAW transaction was inserted for the first time.
|
||||
Inserted,
|
||||
/// Identical canonical RAW transaction content was already durable.
|
||||
AlreadyPresent,
|
||||
/// Normal persistence respected an existing purge tombstone.
|
||||
SkippedPurged,
|
||||
}
|
||||
|
||||
/// Observation disposition produced by one Worker Store persistence attempt.
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub(crate) enum RawTransactionIngestObservationPersistence {
|
||||
/// The deterministic acquisition observation was inserted for the first time.
|
||||
Inserted,
|
||||
/// The same deterministic acquisition observation was already durable.
|
||||
AlreadyPresent,
|
||||
/// No observation was recorded because the canonical entity was intentionally skipped.
|
||||
NotRecorded,
|
||||
}
|
||||
|
||||
/// Classified successful outcome of one atomic Worker Store persistence attempt.
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub(crate) struct RawTransactionIngestPersistenceOutcome {
|
||||
entity: crate::RawTransactionIngestEntityPersistence,
|
||||
observation: crate::RawTransactionIngestObservationPersistence,
|
||||
}
|
||||
|
||||
impl crate::RawTransactionIngestPersistenceOutcome {
|
||||
/// Returns the canonical entity persistence disposition.
|
||||
#[must_use]
|
||||
pub(crate) const fn entity(self) -> crate::RawTransactionIngestEntityPersistence {
|
||||
return self.entity;
|
||||
}
|
||||
|
||||
/// Returns the acquisition-observation persistence disposition.
|
||||
#[must_use]
|
||||
pub(crate) const fn observation(self) -> crate::RawTransactionIngestObservationPersistence {
|
||||
return self.observation;
|
||||
}
|
||||
}
|
||||
|
||||
/// Private backend-neutral Store persistence port used by the Worker and deterministic tests.
|
||||
pub(crate) trait RawTransactionIngestPersistencePort: std::marker::Send + std::marker::Sync {
|
||||
/// Returns whether the persistence target belongs to the requested logical network.
|
||||
fn network_matches(&self, network: &ksp_store_lib::RawNetworkId) -> bool;
|
||||
|
||||
/// Executes one atomic canonical RAW transaction plus observation Store write.
|
||||
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>>;
|
||||
}
|
||||
|
||||
impl crate::RawTransactionIngestPersistencePort for ksp_store_lib::Store {
|
||||
fn network_matches(&self, network: &ksp_store_lib::RawNetworkId) -> bool {
|
||||
return self.runtime_snapshot().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>> {
|
||||
return ksp_store_lib::RawTransactionWrite::persist_raw_transaction_acquisition(self, transaction, observation, mode);
|
||||
}
|
||||
}
|
||||
|
||||
/// Persists one already-canonical Worker acquisition through the private Store port in `Normal` mode.
|
||||
pub(crate) async fn persist_raw_transaction_ingest_acquisition<P>(
|
||||
port: &P,
|
||||
acquisition: ksp_raw_transaction_lib::RawTransactionAcquisition,
|
||||
) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome>
|
||||
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 (transaction, observation) = acquisition.into_parts();
|
||||
let result = port.persist_acquisition(transaction, observation, ksp_store_lib::RawTransactionAcquisitionMode::Normal).await;
|
||||
return match result {
|
||||
std::result::Result::Ok(outcome) => map_store_outcome(outcome),
|
||||
std::result::Result::Err(error) => map_store_error(error),
|
||||
};
|
||||
}
|
||||
|
||||
fn map_store_error(error: ksp_core_lib::Error) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome> {
|
||||
if error.code() == ksp_store_lib::ERROR_CODE_RAW_CONFLICT {
|
||||
return std::result::Result::Err(crate::content_conflict_error());
|
||||
}
|
||||
return std::result::Result::Err(crate::store_error(error.code()));
|
||||
}
|
||||
|
||||
fn map_store_outcome(outcome: ksp_store_lib::RawAcquisitionWriteOutcome) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome> {
|
||||
let entity = outcome.entity();
|
||||
let observation = outcome.observation();
|
||||
if entity == ksp_store_lib::RawEntityWriteOutcome::Inserted && observation == ksp_store_lib::RawObservationWriteOutcome::Inserted {
|
||||
return std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
|
||||
entity: crate::RawTransactionIngestEntityPersistence::Inserted,
|
||||
observation: crate::RawTransactionIngestObservationPersistence::Inserted,
|
||||
});
|
||||
}
|
||||
if entity == ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent && observation == ksp_store_lib::RawObservationWriteOutcome::Inserted {
|
||||
return std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
|
||||
entity: crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
|
||||
observation: crate::RawTransactionIngestObservationPersistence::Inserted,
|
||||
});
|
||||
}
|
||||
if entity == ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent && observation == ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent {
|
||||
return std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
|
||||
entity: crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
|
||||
observation: crate::RawTransactionIngestObservationPersistence::AlreadyPresent,
|
||||
});
|
||||
}
|
||||
if entity == ksp_store_lib::RawEntityWriteOutcome::SkippedPurged && observation == ksp_store_lib::RawObservationWriteOutcome::NotRecorded {
|
||||
return std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
|
||||
entity: crate::RawTransactionIngestEntityPersistence::SkippedPurged,
|
||||
observation: crate::RawTransactionIngestObservationPersistence::NotRecorded,
|
||||
});
|
||||
}
|
||||
return std::result::Result::Err(crate::runtime_error("persistence.store_outcome_invalid"));
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "../unit_tests/persistence.rs"]
|
||||
mod tests;
|
||||
@@ -1,5 +1,8 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs
|
||||
// version: 3
|
||||
// version: 4
|
||||
|
||||
type PersistencePort = std::sync::Arc<dyn crate::RawTransactionIngestPersistencePort + 'static>;
|
||||
type PersistenceTasks = tokio::task::JoinSet<ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome>>;
|
||||
|
||||
/// Runtime-neutral boxed future resolving after one RAW transaction ingest Worker has fully reached a terminal lifecycle state.
|
||||
pub type RawTransactionIngestTerminalFuture<'a> =
|
||||
@@ -85,31 +88,74 @@ fn current_runtime_handle() -> ksp_core_lib::Result<tokio::runtime::Handle> {
|
||||
};
|
||||
}
|
||||
|
||||
async fn drain_admission(settings: &crate::RawTransactionIngestSettings, admission: &mut crate::RawTransactionAdmission) -> bool {
|
||||
let mut clean = true;
|
||||
async fn drain_children(children: &mut tokio::task::JoinSet<()>) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
let mut fault = std::option::Option::None;
|
||||
while let std::option::Option::Some(joined) = children.join_next().await {
|
||||
if joined.is_err() && fault.is_none() {
|
||||
fault = std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
||||
}
|
||||
}
|
||||
return fault;
|
||||
}
|
||||
|
||||
async fn drain_persistence(persistence: &mut PersistenceTasks) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
let mut fault = std::option::Option::None;
|
||||
while let std::option::Option::Some(joined) = persistence.join_next().await {
|
||||
let current = persistence_fault(joined);
|
||||
if fault.is_none() {
|
||||
fault = current;
|
||||
}
|
||||
}
|
||||
return fault;
|
||||
}
|
||||
|
||||
async fn drain_admission_and_persistence(
|
||||
settings: &crate::RawTransactionIngestSettings,
|
||||
admission: &mut crate::RawTransactionAdmission,
|
||||
persistence: &mut PersistenceTasks,
|
||||
port: &std::option::Option<PersistencePort>,
|
||||
) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
let mut fault = std::option::Option::None;
|
||||
admission.close();
|
||||
loop {
|
||||
while persistence.len() >= settings.persistence_concurrency() {
|
||||
let joined = persistence.join_next().await;
|
||||
let current = match joined {
|
||||
std::option::Option::Some(value) => persistence_fault(value),
|
||||
std::option::Option::None => std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID),
|
||||
};
|
||||
if fault.is_none() {
|
||||
fault = current;
|
||||
}
|
||||
}
|
||||
let received = admission.receive(settings.network()).await;
|
||||
match received {
|
||||
std::result::Result::Ok(std::option::Option::Some(_acquisition)) => {},
|
||||
std::result::Result::Ok(std::option::Option::None) => return clean,
|
||||
std::result::Result::Err(_) => clean = false,
|
||||
std::result::Result::Ok(std::option::Option::Some(acquisition)) => {
|
||||
if !spawn_persistence(persistence, port, acquisition) && fault.is_none() {
|
||||
fault = std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
||||
}
|
||||
},
|
||||
std::result::Result::Ok(std::option::Option::None) => break,
|
||||
std::result::Result::Err(_) => {
|
||||
if fault.is_none() {
|
||||
fault = std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
||||
}
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn drain_children(children: &mut tokio::task::JoinSet<()>) -> bool {
|
||||
let mut clean = true;
|
||||
while let std::option::Option::Some(joined) = children.join_next().await {
|
||||
if joined.is_err() {
|
||||
clean = false;
|
||||
}
|
||||
let persistence_fault = drain_persistence(persistence).await;
|
||||
if fault.is_none() {
|
||||
fault = persistence_fault;
|
||||
}
|
||||
return clean;
|
||||
return fault;
|
||||
}
|
||||
|
||||
fn finish_faulted(lifecycle: &mut ksp_worker_api::WorkerLifecycle, sender: &tokio::sync::watch::Sender<ksp_worker_api::WorkerState>) {
|
||||
if lifecycle.fault(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID).is_err() {
|
||||
fn finish_faulted(
|
||||
lifecycle: &mut ksp_worker_api::WorkerLifecycle,
|
||||
sender: &tokio::sync::watch::Sender<ksp_worker_api::WorkerState>,
|
||||
code: ksp_core_lib::ErrorCode,
|
||||
) {
|
||||
if lifecycle.fault(code).is_err() {
|
||||
sender.send_replace(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID));
|
||||
return;
|
||||
}
|
||||
@@ -131,10 +177,34 @@ fn finish_stopped(lifecycle: &mut ksp_worker_api::WorkerLifecycle, sender: &toki
|
||||
return;
|
||||
}
|
||||
|
||||
fn merge_fault(
|
||||
first: std::option::Option<ksp_core_lib::ErrorCode>,
|
||||
next: std::option::Option<ksp_core_lib::ErrorCode>,
|
||||
) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
if first.is_some() {
|
||||
return first;
|
||||
}
|
||||
return next;
|
||||
}
|
||||
|
||||
fn persistence_fault(
|
||||
joined: std::result::Result<ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome>, tokio::task::JoinError>,
|
||||
) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
return match joined {
|
||||
std::result::Result::Ok(std::result::Result::Ok(outcome)) => {
|
||||
let _entity = outcome.entity();
|
||||
let _observation = outcome.observation();
|
||||
std::option::Option::None
|
||||
},
|
||||
std::result::Result::Ok(std::result::Result::Err(error)) => std::option::Option::Some(error.code()),
|
||||
std::result::Result::Err(_) => std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID),
|
||||
};
|
||||
}
|
||||
|
||||
async fn run_supervisor<Spawner>(
|
||||
settings: crate::RawTransactionIngestSettings,
|
||||
mut lifecycle: ksp_worker_api::WorkerLifecycle,
|
||||
_store_guard: std::option::Option<std::sync::Arc<ksp_store_lib::Store>>,
|
||||
port: std::option::Option<PersistencePort>,
|
||||
mut stop_receiver: tokio::sync::watch::Receiver<bool>,
|
||||
terminal_sender: tokio::sync::watch::Sender<ksp_worker_api::WorkerState>,
|
||||
source_spawner: Spawner,
|
||||
@@ -153,19 +223,40 @@ async fn run_supervisor<Spawner>(
|
||||
}
|
||||
terminal_sender.send_replace(lifecycle.state());
|
||||
let mut children = tokio::task::JoinSet::new();
|
||||
let mut persistence = PersistenceTasks::new();
|
||||
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, stop_receiver.clone(), admission_sender);
|
||||
let supervised_clean = supervise_until_stop(&settings, &mut stop_receiver, &mut children, &mut admission).await;
|
||||
let admission_clean = drain_admission(&settings, &mut admission).await;
|
||||
let children_clean = drain_children(&mut children).await;
|
||||
if !supervised_clean || !admission_clean || !children_clean {
|
||||
finish_faulted(&mut lifecycle, &terminal_sender);
|
||||
return;
|
||||
source_spawner(&mut children, source_stop_receiver, admission_sender);
|
||||
let mut fault = supervise_until_stop(&settings, &source_stop_sender, &mut stop_receiver, &mut children, &mut admission, &mut persistence, &port).await;
|
||||
if fault.is_some() {
|
||||
source_stop_sender.send_replace(true);
|
||||
}
|
||||
let drain_fault = drain_admission_and_persistence(&settings, &mut admission, &mut persistence, &port).await;
|
||||
fault = merge_fault(fault, drain_fault);
|
||||
let child_fault = drain_children(&mut children).await;
|
||||
fault = merge_fault(fault, child_fault);
|
||||
match fault {
|
||||
std::option::Option::Some(code) => finish_faulted(&mut lifecycle, &terminal_sender, code),
|
||||
std::option::Option::None => finish_stopped(&mut lifecycle, &terminal_sender),
|
||||
}
|
||||
finish_stopped(&mut lifecycle, &terminal_sender);
|
||||
return;
|
||||
}
|
||||
|
||||
fn spawn_persistence(
|
||||
persistence: &mut PersistenceTasks,
|
||||
port: &std::option::Option<PersistencePort>,
|
||||
acquisition: ksp_raw_transaction_lib::RawTransactionAcquisition,
|
||||
) -> bool {
|
||||
let port = match port {
|
||||
std::option::Option::Some(value) => std::sync::Arc::clone(value),
|
||||
std::option::Option::None => return false,
|
||||
};
|
||||
let _abort_handle = persistence.spawn(async move {
|
||||
return crate::persist_raw_transaction_ingest_acquisition(port.as_ref(), acquisition).await;
|
||||
});
|
||||
return true;
|
||||
}
|
||||
|
||||
fn start_foundation(
|
||||
settings: crate::RawTransactionIngestSettings,
|
||||
runtime: tokio::runtime::Handle,
|
||||
@@ -174,10 +265,10 @@ fn start_foundation(
|
||||
return start_foundation_with_source_spawner(settings, runtime, store_guard, |_children, _stop_receiver, _admission_sender| {});
|
||||
}
|
||||
|
||||
fn start_foundation_with_source_spawner<Spawner>(
|
||||
fn start_foundation_with_port_and_source_spawner<Spawner>(
|
||||
settings: crate::RawTransactionIngestSettings,
|
||||
runtime: tokio::runtime::Handle,
|
||||
store_guard: std::option::Option<std::sync::Arc<ksp_store_lib::Store>>,
|
||||
port: std::option::Option<PersistencePort>,
|
||||
source_spawner: Spawner,
|
||||
) -> ksp_core_lib::Result<crate::RawTransactionIngestHandle>
|
||||
where
|
||||
@@ -197,25 +288,50 @@ where
|
||||
let (stop_sender, stop_receiver) = tokio::sync::watch::channel(false);
|
||||
let (terminal_sender, terminal_receiver) = tokio::sync::watch::channel(lifecycle.state());
|
||||
let handle = crate::RawTransactionIngestHandle { stop_sender, stop_token, terminal_receiver };
|
||||
std::mem::drop(runtime.spawn(run_supervisor(settings, lifecycle, store_guard, stop_receiver, terminal_sender, source_spawner)));
|
||||
std::mem::drop(runtime.spawn(run_supervisor(settings, lifecycle, port, stop_receiver, terminal_sender, source_spawner)));
|
||||
return std::result::Result::Ok(handle);
|
||||
}
|
||||
|
||||
fn start_foundation_with_source_spawner<Spawner>(
|
||||
settings: crate::RawTransactionIngestSettings,
|
||||
runtime: tokio::runtime::Handle,
|
||||
store_guard: std::option::Option<std::sync::Arc<ksp_store_lib::Store>>,
|
||||
source_spawner: Spawner,
|
||||
) -> ksp_core_lib::Result<crate::RawTransactionIngestHandle>
|
||||
where
|
||||
Spawner: FnOnce(&mut tokio::task::JoinSet<()>, tokio::sync::watch::Receiver<bool>, tokio::sync::mpsc::Sender<crate::RawTransactionIngress>)
|
||||
+ std::marker::Send
|
||||
+ 'static,
|
||||
{
|
||||
let port = match store_guard {
|
||||
std::option::Option::Some(store) => {
|
||||
let port: PersistencePort = store;
|
||||
std::option::Option::Some(port)
|
||||
},
|
||||
std::option::Option::None => std::option::Option::None,
|
||||
};
|
||||
return start_foundation_with_port_and_source_spawner(settings, runtime, port, source_spawner);
|
||||
}
|
||||
|
||||
async fn supervise_until_stop(
|
||||
settings: &crate::RawTransactionIngestSettings,
|
||||
source_stop_sender: &tokio::sync::watch::Sender<bool>,
|
||||
stop_receiver: &mut tokio::sync::watch::Receiver<bool>,
|
||||
children: &mut tokio::task::JoinSet<()>,
|
||||
admission: &mut crate::RawTransactionAdmission,
|
||||
) -> bool {
|
||||
persistence: &mut PersistenceTasks,
|
||||
port: &std::option::Option<PersistencePort>,
|
||||
) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
let mut admission_open = true;
|
||||
let mut clean = true;
|
||||
loop {
|
||||
if *stop_receiver.borrow() {
|
||||
source_stop_sender.send_replace(true);
|
||||
break;
|
||||
}
|
||||
if children.is_empty() && !admission_open {
|
||||
if children.is_empty() && !admission_open && persistence.is_empty() {
|
||||
let changed = stop_receiver.changed().await;
|
||||
if changed.is_err() || *stop_receiver.borrow() {
|
||||
source_stop_sender.send_replace(true);
|
||||
break;
|
||||
}
|
||||
continue;
|
||||
@@ -224,24 +340,43 @@ async fn supervise_until_stop(
|
||||
biased;
|
||||
changed = stop_receiver.changed() => {
|
||||
if changed.is_err() || *stop_receiver.borrow() {
|
||||
source_stop_sender.send_replace(true);
|
||||
break;
|
||||
}
|
||||
}
|
||||
joined = children.join_next(), if !children.is_empty() => {
|
||||
if let std::option::Option::Some(std::result::Result::Err(_)) = joined {
|
||||
clean = false;
|
||||
source_stop_sender.send_replace(true);
|
||||
return std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
||||
}
|
||||
}
|
||||
received = admission.receive(settings.network()), if admission_open => {
|
||||
joined = persistence.join_next(), if !persistence.is_empty() => {
|
||||
if let std::option::Option::Some(value) = joined {
|
||||
let fault = persistence_fault(value);
|
||||
if fault.is_some() {
|
||||
source_stop_sender.send_replace(true);
|
||||
return fault;
|
||||
}
|
||||
}
|
||||
}
|
||||
received = admission.receive(settings.network()), if admission_open && persistence.len() < settings.persistence_concurrency() => {
|
||||
match received {
|
||||
std::result::Result::Ok(std::option::Option::Some(_acquisition)) => {},
|
||||
std::result::Result::Ok(std::option::Option::Some(acquisition)) => {
|
||||
if !spawn_persistence(persistence, port, acquisition) {
|
||||
source_stop_sender.send_replace(true);
|
||||
return std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
||||
}
|
||||
},
|
||||
std::result::Result::Ok(std::option::Option::None) => admission_open = false,
|
||||
std::result::Result::Err(_) => clean = false,
|
||||
std::result::Result::Err(_) => {
|
||||
source_stop_sender.send_replace(true);
|
||||
return std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return clean;
|
||||
return std::option::Option::None;
|
||||
}
|
||||
|
||||
fn validate_store_network(settings: &crate::RawTransactionIngestSettings, store_network: &ksp_store_lib::RawNetworkId) -> ksp_core_lib::Result<()> {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs
|
||||
// version: 5
|
||||
// version: 6
|
||||
|
||||
//! Dependency firewall canaries for the RAW transaction ingest Worker foundation.
|
||||
|
||||
@@ -52,36 +52,41 @@ fn pre_002_manifest_keeps_forbidden_layers_and_live_sources_out() {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pre_006_source_surface_opens_bounded_admission_and_common_raw_without_persistence_or_live_sources() {
|
||||
fn pre_007_source_surface_adds_private_normal_store_persistence_without_backend_or_live_source() {
|
||||
let root = include_str!("../src/lib.rs");
|
||||
let runtime = include_str!("../src/runtime.rs");
|
||||
let admission = include_str!("../src/admission.rs");
|
||||
let persistence = include_str!("../src/persistence.rs");
|
||||
for required in [
|
||||
"tokio::task::JoinSet",
|
||||
"tokio::sync::mpsc::channel",
|
||||
"canonicalize_raw_transaction",
|
||||
"assemble_raw_transaction_acquisition",
|
||||
"ksp.raw_transaction_ingest.observation.v1",
|
||||
"RawTransactionIngress",
|
||||
"RawTransactionIngestPersistencePort",
|
||||
"persist_raw_transaction_acquisition",
|
||||
"RawTransactionAcquisitionMode::Normal",
|
||||
"persistence_concurrency()",
|
||||
"ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT",
|
||||
"ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED",
|
||||
] {
|
||||
assert!(
|
||||
root.contains(required) || runtime.contains(required) || admission.contains(required),
|
||||
"required pre.006 bounded admission/common RAW contract missing: {required}"
|
||||
root.contains(required) || runtime.contains(required) || admission.contains(required) || persistence.contains(required),
|
||||
"required pre.007 persistence contract missing: {required}"
|
||||
);
|
||||
}
|
||||
for forbidden in [
|
||||
"pub use self::admission::RawTransactionAdmission",
|
||||
"pub use self::admission::RawTransactionIngress",
|
||||
"pub use self::persistence::RawTransactionIngestPersistencePort",
|
||||
"unbounded_channel",
|
||||
"persist_raw_transaction",
|
||||
"RawTransactionAcquisitionMode",
|
||||
"ksp_store_postgres_lib::",
|
||||
"ksp_onchain_transport_lib::",
|
||||
"RawTransactionIngestSnapshot",
|
||||
"WorkerSnapshotSource",
|
||||
"ForceRehydrate",
|
||||
] {
|
||||
assert!(
|
||||
!root.contains(forbidden) && !runtime.contains(forbidden) && !admission.contains(forbidden),
|
||||
"pre.006 opened later runtime scope too early: {forbidden}"
|
||||
!root.contains(forbidden) && !runtime.contains(forbidden) && !admission.contains(forbidden) && !persistence.contains(forbidden),
|
||||
"pre.007 crossed a forbidden runtime boundary: {forbidden}"
|
||||
);
|
||||
}
|
||||
return;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/tests/public_api.rs
|
||||
// version: 5
|
||||
// version: 6
|
||||
|
||||
//! External public-surface proofs for the RAW transaction ingest Worker identity and settings foundation.
|
||||
|
||||
@@ -91,3 +91,12 @@ fn pre_004_start_handle_and_terminal_future_are_consumable_without_public_join_h
|
||||
assert!(!root.contains("JoinSet"));
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pre_007_store_and_content_conflict_error_codes_are_public_and_stable() {
|
||||
assert_eq!(ksp_worker_raw_transaction_ingest_lib::ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT.domain(), "worker_raw_transaction_ingest");
|
||||
assert_eq!(ksp_worker_raw_transaction_ingest_lib::ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT.code(), "content_conflict");
|
||||
assert_eq!(ksp_worker_raw_transaction_ingest_lib::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED.domain(), "worker_raw_transaction_ingest");
|
||||
assert_eq!(ksp_worker_raw_transaction_ingest_lib::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED.code(), "store_failed");
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,253 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/persistence.rs
|
||||
// version: 1
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
enum PortResponse {
|
||||
Outcome(ksp_store_lib::RawEntityWriteOutcome, ksp_store_lib::RawObservationWriteOutcome),
|
||||
Conflict,
|
||||
StoreFailure,
|
||||
}
|
||||
|
||||
struct FakePersistencePort {
|
||||
calls: std::sync::atomic::AtomicUsize,
|
||||
network: ksp_store_lib::RawNetworkId,
|
||||
normal_mode_seen: std::sync::atomic::AtomicBool,
|
||||
response: PortResponse,
|
||||
}
|
||||
|
||||
impl FakePersistencePort {
|
||||
fn new(network: ksp_store_lib::RawNetworkId, response: PortResponse) -> Self {
|
||||
return Self {
|
||||
calls: std::sync::atomic::AtomicUsize::new(0),
|
||||
network,
|
||||
normal_mode_seen: std::sync::atomic::AtomicBool::new(false),
|
||||
response,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
impl crate::RawTransactionIngestPersistencePort for FakePersistencePort {
|
||||
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.calls.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
|
||||
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 {
|
||||
return match response {
|
||||
PortResponse::Outcome(entity, observation) => std::result::Result::Ok(ksp_store_lib::RawAcquisitionWriteOutcome::new(entity, 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<ksp_raw_transaction_lib::RawTransactionAcquisition> {
|
||||
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,
|
||||
};
|
||||
let signature = ksp_store_lib::RawTransactionSignature::new([signature_byte; 64]);
|
||||
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64(
|
||||
network,
|
||||
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 transaction = match ksp_raw_transaction_lib::canonicalize_raw_transaction(material) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::option::Option::None,
|
||||
};
|
||||
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("persistence") {
|
||||
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(ksp_raw_transaction_lib::assemble_raw_transaction_acquisition(
|
||||
transaction,
|
||||
ksp_store_lib::RawObservationKey::new([signature_byte; 32]),
|
||||
provenance,
|
||||
));
|
||||
}
|
||||
|
||||
fn network(name: &str) -> std::option::Option<ksp_store_lib::RawNetworkId> {
|
||||
return match ksp_store_lib::RawNetworkId::new(name) {
|
||||
std::result::Result::Ok(value) => std::option::Option::Some(value),
|
||||
std::result::Result::Err(_) => std::option::Option::None,
|
||||
};
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn pre_007_normal_mode_classifies_new_idempotent_and_purged_outcomes() {
|
||||
let cases = [
|
||||
(
|
||||
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::Inserted),
|
||||
crate::RawTransactionIngestEntityPersistence::Inserted,
|
||||
crate::RawTransactionIngestObservationPersistence::Inserted,
|
||||
),
|
||||
(
|
||||
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent, ksp_store_lib::RawObservationWriteOutcome::Inserted),
|
||||
crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
|
||||
crate::RawTransactionIngestObservationPersistence::Inserted,
|
||||
),
|
||||
(
|
||||
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent, ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent),
|
||||
crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
|
||||
crate::RawTransactionIngestObservationPersistence::AlreadyPresent,
|
||||
),
|
||||
(
|
||||
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::SkippedPurged, ksp_store_lib::RawObservationWriteOutcome::NotRecorded),
|
||||
crate::RawTransactionIngestEntityPersistence::SkippedPurged,
|
||||
crate::RawTransactionIngestObservationPersistence::NotRecorded,
|
||||
),
|
||||
];
|
||||
for (index, (response, expected_entity, expected_observation)) in cases.into_iter().enumerate() {
|
||||
let network = match network("mainnet") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let port = FakePersistencePort::new(network, response);
|
||||
let acquisition = match acquisition("mainnet", index as u8 + 1) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let result = crate::persist_raw_transaction_ingest_acquisition(&port, acquisition).await;
|
||||
assert!(result.is_ok());
|
||||
let outcome = match result {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return,
|
||||
};
|
||||
assert_eq!(outcome.entity(), expected_entity);
|
||||
assert_eq!(outcome.observation(), expected_observation);
|
||||
assert!(port.normal_mode_seen.load(std::sync::atomic::Ordering::Acquire));
|
||||
assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 1);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn pre_007_content_conflict_maps_to_terminal_worker_code_without_remote_text() {
|
||||
let network = match network("mainnet") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let port = FakePersistencePort::new(network, PortResponse::Conflict);
|
||||
let acquisition = match acquisition("mainnet", 9) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let result = crate::persist_raw_transaction_ingest_acquisition(&port, acquisition).await;
|
||||
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_CONTENT_CONFLICT);
|
||||
assert!(error.context().is_empty());
|
||||
assert!(!std::format!("{error:?}").contains("synthetic conflict"));
|
||||
return;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn pre_007_store_failure_keeps_only_safe_lower_error_code() {
|
||||
let network = match network("mainnet") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let port = FakePersistencePort::new(network, PortResponse::StoreFailure);
|
||||
let acquisition = match acquisition("mainnet", 10) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let result = crate::persist_raw_transaction_ingest_acquisition(&port, acquisition).await;
|
||||
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_STORE_FAILED);
|
||||
assert_eq!(error.context().len(), 2);
|
||||
assert_eq!(error.context()[0].value(), "synthetic_store");
|
||||
assert_eq!(error.context()[1].value(), "write_failed");
|
||||
assert!(!std::format!("{error:?}").contains("synthetic remote failure text"));
|
||||
return;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn pre_007_rehydrated_or_impossible_store_outcomes_are_runtime_invalid() {
|
||||
let cases = [
|
||||
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::Rehydrated, ksp_store_lib::RawObservationWriteOutcome::Inserted),
|
||||
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent),
|
||||
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::SkippedPurged, ksp_store_lib::RawObservationWriteOutcome::Inserted),
|
||||
];
|
||||
for (index, response) in cases.into_iter().enumerate() {
|
||||
let network = match network("mainnet") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let port = FakePersistencePort::new(network, response);
|
||||
let acquisition = match acquisition("mainnet", index as u8 + 20) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let result = crate::persist_raw_transaction_ingest_acquisition(&port, acquisition).await;
|
||||
let error = match result {
|
||||
std::result::Result::Ok(_) => continue,
|
||||
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(), "persistence.store_outcome_invalid");
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn pre_007_store_network_mismatch_is_rejected_before_write() {
|
||||
let network = match network("devnet") {
|
||||
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),
|
||||
);
|
||||
let acquisition = match acquisition("mainnet", 30) {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let result = crate::persist_raw_transaction_ingest_acquisition(&port, acquisition).await;
|
||||
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()[0].value(), "persistence.store_network_mismatch");
|
||||
assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 0);
|
||||
return;
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs
|
||||
// version: 3
|
||||
// version: 4
|
||||
|
||||
struct ActiveTaskGuard {
|
||||
active: std::sync::Arc<std::sync::atomic::AtomicUsize>,
|
||||
@@ -251,3 +251,284 @@ async fn pre_005_supervisor_reaps_completed_source_task_and_still_joins_live_chi
|
||||
assert_eq!(active.load(std::sync::atomic::Ordering::Acquire), 0);
|
||||
return;
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
enum RuntimePortResponse {
|
||||
Conflict,
|
||||
StoreFailure,
|
||||
SuccessBlocked,
|
||||
}
|
||||
|
||||
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::SuccessBlocked) {
|
||||
let current = self.active.fetch_add(1, std::sync::atomic::Ordering::AcqRel) + 1;
|
||||
self.update_max_active(current);
|
||||
while !self.released.load(std::sync::atomic::Ordering::Acquire) {
|
||||
self.notify.notified().await;
|
||||
}
|
||||
self.active.fetch_sub(1, std::sync::atomic::Ordering::AcqRel);
|
||||
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 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> {
|
||||
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, 8, concurrency, std::time::Duration::from_secs(5)) {
|
||||
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,
|
||||
};
|
||||
if admission_sender.send(ingress).await.is_err() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
return;
|
||||
});
|
||||
},
|
||||
) {
|
||||
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,
|
||||
};
|
||||
let _sent = admission_sender.send(ingress).await;
|
||||
return;
|
||||
});
|
||||
},
|
||||
) {
|
||||
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_CONTENT_CONFLICT));
|
||||
assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 1);
|
||||
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,
|
||||
};
|
||||
let _sent = admission_sender.send(ingress).await;
|
||||
return;
|
||||
});
|
||||
},
|
||||
) {
|
||||
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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user