v0.3.11-pre.003
This commit is contained in:
12
crates/ksp-worker-raw-transaction-ingest-lib/src/error.rs
Normal file
12
crates/ksp-worker-raw-transaction-ingest-lib/src/error.rs
Normal file
@@ -0,0 +1,12 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/error.rs
|
||||
// version: 2
|
||||
|
||||
/// 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");
|
||||
|
||||
/// Creates one settings-domain error carrying only the stable invalid field name.
|
||||
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);
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/identity.rs
|
||||
// version: 1
|
||||
|
||||
/// Stable Worker kind code used by the continuous RAW transaction ingest vertical.
|
||||
pub const RAW_TRANSACTION_INGEST_WORKER_KIND_CODE: &str = "raw_transaction_ingest";
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs
|
||||
// version: 1
|
||||
// version: 2
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -7,6 +7,38 @@
|
||||
|
||||
//! Source-neutral runtime foundation for continuous KSP RAW transaction ingestion.
|
||||
//!
|
||||
//! The crate is intentionally limited to its dependency skeleton in this tranche.
|
||||
//! Worker identity, settings, lifecycle, tasks, admission, persistence and snapshots
|
||||
//! are introduced only by their dedicated prereleases.
|
||||
//! This tranche owns only the concrete Worker family identity and validated technical
|
||||
//! settings. Lifecycle, tasks, admission, persistence and snapshots are introduced only
|
||||
//! by their dedicated prereleases; no live source or Transport dependency exists here.
|
||||
|
||||
mod error;
|
||||
mod identity;
|
||||
mod settings;
|
||||
|
||||
/// 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;
|
||||
/// Stable Worker kind code used by the continuous RAW transaction ingest vertical.
|
||||
pub use self::identity::RAW_TRANSACTION_INGEST_WORKER_KIND_CODE;
|
||||
/// Default bounded admission queue capacity for one RAW transaction ingest Worker.
|
||||
pub use self::settings::DEFAULT_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY;
|
||||
/// Default number of concurrent Store persistence operations for one RAW transaction ingest Worker.
|
||||
pub use self::settings::DEFAULT_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY;
|
||||
/// Default cooperative shutdown drain deadline for one RAW transaction ingest Worker.
|
||||
pub use self::settings::DEFAULT_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT;
|
||||
/// Maximum bounded admission queue capacity for one RAW transaction ingest Worker.
|
||||
pub use self::settings::MAX_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY;
|
||||
/// Maximum number of concurrent Store persistence operations for one RAW transaction ingest Worker.
|
||||
pub use self::settings::MAX_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY;
|
||||
/// Maximum cooperative shutdown drain deadline for one RAW transaction ingest Worker.
|
||||
pub use self::settings::MAX_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT;
|
||||
/// Minimum bounded admission queue capacity for one RAW transaction ingest Worker.
|
||||
pub use self::settings::MIN_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY;
|
||||
/// Minimum number of concurrent Store persistence operations for one RAW transaction ingest Worker.
|
||||
pub use self::settings::MIN_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY;
|
||||
/// Minimum cooperative shutdown drain deadline for one RAW transaction ingest Worker.
|
||||
pub use self::settings::MIN_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT;
|
||||
/// Validated source-neutral runtime settings for one continuous RAW transaction ingest Worker.
|
||||
pub use self::settings::RawTransactionIngestSettings;
|
||||
|
||||
/// Creates one settings-domain error without copying caller-supplied values into diagnostics.
|
||||
pub(crate) use self::error::settings_error;
|
||||
|
||||
127
crates/ksp-worker-raw-transaction-ingest-lib/src/settings.rs
Normal file
127
crates/ksp-worker-raw-transaction-ingest-lib/src/settings.rs
Normal file
@@ -0,0 +1,127 @@
|
||||
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/settings.rs
|
||||
// version: 3
|
||||
|
||||
/// Default bounded admission queue capacity for one RAW transaction ingest Worker.
|
||||
pub const DEFAULT_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY: usize = 256;
|
||||
/// Default number of concurrent Store persistence operations for one RAW transaction ingest Worker.
|
||||
pub const DEFAULT_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY: usize = 8;
|
||||
/// Default cooperative shutdown drain deadline for one RAW transaction ingest Worker.
|
||||
pub const DEFAULT_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
|
||||
/// Maximum bounded admission queue capacity for one RAW transaction ingest Worker.
|
||||
pub const MAX_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY: usize = 65_536;
|
||||
/// Maximum number of concurrent Store persistence operations for one RAW transaction ingest Worker.
|
||||
pub const MAX_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY: usize = 64;
|
||||
/// Maximum cooperative shutdown drain deadline for one RAW transaction ingest Worker.
|
||||
pub const MAX_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
|
||||
/// Minimum bounded admission queue capacity for one RAW transaction ingest Worker.
|
||||
pub const MIN_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY: usize = 1;
|
||||
/// Minimum number of concurrent Store persistence operations for one RAW transaction ingest Worker.
|
||||
pub const MIN_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY: usize = 1;
|
||||
/// Minimum cooperative shutdown drain deadline for one RAW transaction ingest Worker.
|
||||
pub const MIN_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(100);
|
||||
|
||||
/// Validated source-neutral runtime settings for one continuous RAW transaction ingest Worker.
|
||||
#[derive(Clone, Eq, PartialEq)]
|
||||
pub struct RawTransactionIngestSettings {
|
||||
network: ksp_store_lib::RawNetworkId,
|
||||
worker_id: ksp_worker_api::WorkerId,
|
||||
admission_queue_capacity: usize,
|
||||
persistence_concurrency: usize,
|
||||
shutdown_drain_timeout: std::time::Duration,
|
||||
}
|
||||
|
||||
impl crate::RawTransactionIngestSettings {
|
||||
/// Creates settings from validated network/Worker identities and explicit bounded runtime limits.
|
||||
pub fn new(
|
||||
network: ksp_store_lib::RawNetworkId,
|
||||
worker_id: ksp_worker_api::WorkerId,
|
||||
admission_queue_capacity: usize,
|
||||
persistence_concurrency: usize,
|
||||
shutdown_drain_timeout: std::time::Duration,
|
||||
) -> ksp_core_lib::Result<Self> {
|
||||
validate_inclusive_usize(
|
||||
admission_queue_capacity,
|
||||
crate::MIN_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY,
|
||||
crate::MAX_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY,
|
||||
"admission_queue_capacity",
|
||||
)?;
|
||||
validate_inclusive_usize(
|
||||
persistence_concurrency,
|
||||
crate::MIN_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY,
|
||||
crate::MAX_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY,
|
||||
"persistence_concurrency",
|
||||
)?;
|
||||
if !(crate::MIN_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT..=crate::MAX_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT)
|
||||
.contains(&shutdown_drain_timeout)
|
||||
{
|
||||
return std::result::Result::Err(crate::settings_error("shutdown_drain_timeout"));
|
||||
}
|
||||
return std::result::Result::Ok(Self { network, worker_id, admission_queue_capacity, persistence_concurrency, shutdown_drain_timeout });
|
||||
}
|
||||
|
||||
/// Creates settings using the stable V1 runtime defaults for one validated network and Worker identity.
|
||||
#[must_use]
|
||||
pub fn with_defaults(network: ksp_store_lib::RawNetworkId, worker_id: ksp_worker_api::WorkerId) -> Self {
|
||||
return Self {
|
||||
network,
|
||||
worker_id,
|
||||
admission_queue_capacity: crate::DEFAULT_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY,
|
||||
persistence_concurrency: crate::DEFAULT_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY,
|
||||
shutdown_drain_timeout: crate::DEFAULT_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT,
|
||||
};
|
||||
}
|
||||
|
||||
/// Returns the Store-compatible logical network binding.
|
||||
#[must_use]
|
||||
pub const fn network(&self) -> &ksp_store_lib::RawNetworkId {
|
||||
return &self.network;
|
||||
}
|
||||
|
||||
/// Returns the validated logical identity of this Worker instance.
|
||||
#[must_use]
|
||||
pub const fn worker_id(&self) -> &ksp_worker_api::WorkerId {
|
||||
return &self.worker_id;
|
||||
}
|
||||
|
||||
/// Returns the bounded central admission queue capacity.
|
||||
#[must_use]
|
||||
pub const fn admission_queue_capacity(&self) -> usize {
|
||||
return self.admission_queue_capacity;
|
||||
}
|
||||
|
||||
/// Returns the bounded maximum number of concurrent Store persistence operations.
|
||||
#[must_use]
|
||||
pub const fn persistence_concurrency(&self) -> usize {
|
||||
return self.persistence_concurrency;
|
||||
}
|
||||
|
||||
/// Returns the bounded cooperative shutdown drain deadline.
|
||||
#[must_use]
|
||||
pub const fn shutdown_drain_timeout(&self) -> std::time::Duration {
|
||||
return self.shutdown_drain_timeout;
|
||||
}
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for crate::RawTransactionIngestSettings {
|
||||
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
return formatter
|
||||
.debug_struct("RawTransactionIngestSettings")
|
||||
.field("network", &self.network)
|
||||
.field("worker_id", &self.worker_id)
|
||||
.field("admission_queue_capacity", &self.admission_queue_capacity)
|
||||
.field("persistence_concurrency", &self.persistence_concurrency)
|
||||
.field("shutdown_drain_timeout", &self.shutdown_drain_timeout)
|
||||
.finish();
|
||||
}
|
||||
}
|
||||
|
||||
fn validate_inclusive_usize(value: usize, minimum: usize, maximum: usize, field: &'static str) -> ksp_core_lib::Result<()> {
|
||||
if !(minimum..=maximum).contains(&value) {
|
||||
return std::result::Result::Err(crate::settings_error(field));
|
||||
}
|
||||
return std::result::Result::Ok(());
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "../unit_tests/settings.rs"]
|
||||
mod tests;
|
||||
Reference in New Issue
Block a user