v0.3.11-pre.008

This commit is contained in:
2026-09-08 12:28:51 +02:00
parent ed6c0ac10c
commit 22c88c449d
13 changed files with 1241 additions and 88 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs
// version: 4
// version: 5
type PersistencePort = std::sync::Arc<dyn crate::RawTransactionIngestPersistencePort + 'static>;
type PersistenceTasks = tokio::task::JoinSet<ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome>>;
@@ -11,9 +11,9 @@ pub type RawTransactionIngestTerminalFuture<'a> =
/// Cloneable external control handle for one continuous RAW transaction ingest Worker.
#[derive(Clone)]
pub struct RawTransactionIngestHandle {
snapshots: crate::RawTransactionIngestSnapshotSource,
stop_sender: tokio::sync::watch::Sender<bool>,
stop_token: ksp_worker_api::WorkerStopToken,
terminal_receiver: tokio::sync::watch::Receiver<ksp_worker_api::WorkerState>,
}
impl crate::RawTransactionIngestHandle {
@@ -26,22 +26,36 @@ impl crate::RawTransactionIngestHandle {
return self.stop_sender.send(true).is_ok();
}
/// Waits until the private runtime task has published and closed one terminal lifecycle state.
/// Returns an independent latest-value source for concrete RAW transaction ingest Worker snapshots.
#[must_use]
pub fn snapshot_source(&self) -> crate::RawTransactionIngestSnapshotSource {
return self.snapshots.clone();
}
/// Returns an independent source implementing the common [`ksp_worker_api::WorkerSnapshotSource`] projection.
#[must_use]
pub fn worker_snapshot_source(&self) -> crate::RawTransactionIngestSnapshotSource {
return self.snapshots.clone();
}
/// Waits until the private runtime task has published one terminal lifecycle state after draining owned work.
#[must_use]
pub fn wait_terminal(&self) -> crate::RawTransactionIngestTerminalFuture<'_> {
let mut receiver = self.terminal_receiver.clone();
let source = self.snapshots.clone();
return std::boxed::Box::pin(async move {
loop {
let current = *receiver.borrow();
if current.is_terminal() {
let changed = receiver.changed().await;
if changed.is_err() {
return std::result::Result::Ok(current);
}
return std::result::Result::Err(crate::runtime_error("terminal.changed_after_terminal"));
let current = source.current();
let state = current.worker_snapshot().state();
if state.is_terminal() {
return std::result::Result::Ok(state);
}
let changed = receiver.changed().await;
if changed.is_err() {
let observed = current.worker_snapshot().sequence();
let changed = source.wait_for_change(observed).await;
let state = changed.worker_snapshot().state();
if state.is_terminal() {
return std::result::Result::Ok(state);
}
if changed.worker_snapshot().sequence() == observed && source.is_closed() {
return std::result::Result::Err(crate::runtime_error("terminal.closed_before_terminal"));
}
}
@@ -51,11 +65,12 @@ impl crate::RawTransactionIngestHandle {
impl std::fmt::Debug for crate::RawTransactionIngestHandle {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let state = *self.terminal_receiver.borrow();
let snapshot = self.snapshots.current();
return formatter
.debug_struct("RawTransactionIngestHandle")
.field("stop_requested", &self.stop_token.is_stop_requested())
.field("state", &state)
.field("sequence", &snapshot.worker_snapshot().sequence())
.field("state", &snapshot.worker_snapshot().state())
.finish();
}
}
@@ -81,6 +96,21 @@ impl crate::RawTransactionIngestWorker {
}
}
fn begin_stopping(
lifecycle: &mut ksp_worker_api::WorkerLifecycle,
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
admission_queue_depth: usize,
in_flight_persistence: usize,
) -> std::option::Option<ksp_core_lib::ErrorCode> {
if lifecycle.mark_stopping().is_err() {
return std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
}
if let std::result::Result::Err(error) = snapshots.publish_state(lifecycle.state(), admission_queue_depth, in_flight_persistence) {
return std::option::Option::Some(error.code());
}
return std::option::Option::None;
}
fn current_runtime_handle() -> ksp_core_lib::Result<tokio::runtime::Handle> {
return match tokio::runtime::Handle::try_current() {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
@@ -98,10 +128,15 @@ async fn drain_children(children: &mut tokio::task::JoinSet<()>) -> std::option:
return fault;
}
async fn drain_persistence(persistence: &mut PersistenceTasks) -> std::option::Option<ksp_core_lib::ErrorCode> {
async fn drain_persistence(
lifecycle: &ksp_worker_api::WorkerLifecycle,
admission: &crate::RawTransactionAdmission,
persistence: &mut PersistenceTasks,
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
) -> 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);
let current = persistence_completion(joined, lifecycle.state(), admission.queue_depth(), persistence.len(), snapshots);
if fault.is_none() {
fault = current;
}
@@ -111,9 +146,11 @@ async fn drain_persistence(persistence: &mut PersistenceTasks) -> std::option::O
async fn drain_admission_and_persistence(
settings: &crate::RawTransactionIngestSettings,
lifecycle: &ksp_worker_api::WorkerLifecycle,
admission: &mut crate::RawTransactionAdmission,
persistence: &mut PersistenceTasks,
port: &std::option::Option<PersistencePort>,
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
) -> std::option::Option<ksp_core_lib::ErrorCode> {
let mut fault = std::option::Option::None;
admission.close();
@@ -121,7 +158,7 @@ async fn drain_admission_and_persistence(
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::Some(value) => persistence_completion(value, lifecycle.state(), admission.queue_depth(), persistence.len(), snapshots),
std::option::Option::None => std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID),
};
if fault.is_none() {
@@ -131,19 +168,31 @@ async fn drain_admission_and_persistence(
let received = admission.receive(settings.network()).await;
match received {
std::result::Result::Ok(std::option::Option::Some(acquisition)) => {
if !spawn_persistence(persistence, port, acquisition) && fault.is_none() {
let spawned = spawn_persistence(persistence, port, acquisition);
let published = snapshots.record_admission_success(lifecycle.state(), admission.queue_depth(), persistence.len());
if !spawned && fault.is_none() {
fault = std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
}
if let std::result::Result::Err(error) = published {
if fault.is_none() {
fault = std::option::Option::Some(error.code());
}
}
},
std::result::Result::Ok(std::option::Option::None) => break,
std::result::Result::Err(_) => {
if fault.is_none() {
let published = snapshots.record_admission_failure(lifecycle.state(), admission.queue_depth(), persistence.len());
if let std::result::Result::Err(error) = published {
if fault.is_none() {
fault = std::option::Option::Some(error.code());
}
} else if fault.is_none() {
fault = std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
}
},
}
}
let persistence_fault = drain_persistence(persistence).await;
let persistence_fault = drain_persistence(lifecycle, admission, persistence, snapshots).await;
if fault.is_none() {
fault = persistence_fault;
}
@@ -152,28 +201,27 @@ async fn drain_admission_and_persistence(
fn finish_faulted(
lifecycle: &mut ksp_worker_api::WorkerLifecycle,
sender: &tokio::sync::watch::Sender<ksp_worker_api::WorkerState>,
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
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));
snapshots.force_terminal(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID));
return;
}
sender.send_replace(lifecycle.state());
if snapshots.publish_state(lifecycle.state(), 0, 0).is_err() {
snapshots.force_terminal(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_COUNTER_EXHAUSTED));
}
return;
}
fn finish_stopped(lifecycle: &mut ksp_worker_api::WorkerLifecycle, sender: &tokio::sync::watch::Sender<ksp_worker_api::WorkerState>) {
if lifecycle.mark_stopping().is_err() {
sender.send_replace(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID));
return;
}
sender.send_replace(lifecycle.state());
fn finish_stopped(lifecycle: &mut ksp_worker_api::WorkerLifecycle, snapshots: &mut crate::RawTransactionIngestSnapshotPublisher) {
if lifecycle.mark_stopped().is_err() {
sender.send_replace(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID));
snapshots.force_terminal(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID));
return;
}
sender.send_replace(lifecycle.state());
if snapshots.publish_state(lifecycle.state(), 0, 0).is_err() {
snapshots.force_terminal(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_COUNTER_EXHAUSTED));
}
return;
}
@@ -187,16 +235,28 @@ fn merge_fault(
return next;
}
fn persistence_fault(
fn persistence_completion(
joined: std::result::Result<ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome>, tokio::task::JoinError>,
state: ksp_worker_api::WorkerState,
admission_queue_depth: usize,
in_flight_persistence: usize,
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
) -> 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
match snapshots.record_persistence_success(state, admission_queue_depth, in_flight_persistence, outcome) {
std::result::Result::Ok(()) => std::option::Option::None,
std::result::Result::Err(error) => std::option::Option::Some(error.code()),
}
},
std::result::Result::Ok(std::result::Result::Err(error)) => {
let code = error.code();
let published = snapshots.record_persistence_fault(state, admission_queue_depth, in_flight_persistence, code);
match published {
std::result::Result::Ok(()) => std::option::Option::Some(code),
std::result::Result::Err(snapshot_error) => std::option::Option::Some(snapshot_error.code()),
}
},
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),
};
}
@@ -206,7 +266,7 @@ async fn run_supervisor<Spawner>(
mut lifecycle: ksp_worker_api::WorkerLifecycle,
port: std::option::Option<PersistencePort>,
mut stop_receiver: tokio::sync::watch::Receiver<bool>,
terminal_sender: tokio::sync::watch::Sender<ksp_worker_api::WorkerState>,
mut snapshots: crate::RawTransactionIngestSnapshotPublisher,
source_spawner: Spawner,
) where
Spawner: FnOnce(&mut tokio::task::JoinSet<()>, tokio::sync::watch::Receiver<bool>, tokio::sync::mpsc::Sender<crate::RawTransactionIngress>)
@@ -214,30 +274,50 @@ async fn run_supervisor<Spawner>(
+ 'static,
{
if *stop_receiver.borrow() {
finish_stopped(&mut lifecycle, &terminal_sender);
let stopping_fault = begin_stopping(&mut lifecycle, &mut snapshots, 0, 0);
match stopping_fault {
std::option::Option::Some(code) => finish_faulted(&mut lifecycle, &mut snapshots, code),
std::option::Option::None => finish_stopped(&mut lifecycle, &mut snapshots),
}
return;
}
if lifecycle.mark_running().is_err() {
terminal_sender.send_replace(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID));
snapshots.force_terminal(ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID));
return;
}
if let std::result::Result::Err(error) = snapshots.publish_state(lifecycle.state(), 0, 0) {
finish_faulted(&mut lifecycle, &mut snapshots, error.code());
return;
}
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, 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;
let mut fault = supervise_until_stop(
&settings,
&lifecycle,
&source_stop_sender,
&mut stop_receiver,
&mut children,
&mut admission,
&mut persistence,
&port,
&mut snapshots,
)
.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;
let stopping_fault = begin_stopping(&mut lifecycle, &mut snapshots, admission.queue_depth(), persistence.len());
fault = merge_fault(fault, stopping_fault);
let drain_fault = drain_admission_and_persistence(&settings, &lifecycle, &mut admission, &mut persistence, &port, &mut snapshots).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),
std::option::Option::Some(code) => finish_faulted(&mut lifecycle, &mut snapshots, code),
std::option::Option::None => finish_stopped(&mut lifecycle, &mut snapshots),
}
return;
}
@@ -286,9 +366,9 @@ where
}
let stop_token = ksp_worker_api::WorkerStopToken::new();
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, port, stop_receiver, terminal_sender, source_spawner)));
let (snapshots, snapshot_source) = crate::RawTransactionIngestSnapshotPublisher::new(&settings, &lifecycle);
let handle = crate::RawTransactionIngestHandle { snapshots: snapshot_source, stop_sender, stop_token };
std::mem::drop(runtime.spawn(run_supervisor(settings, lifecycle, port, stop_receiver, snapshots, source_spawner)));
return std::result::Result::Ok(handle);
}
@@ -315,12 +395,14 @@ where
async fn supervise_until_stop(
settings: &crate::RawTransactionIngestSettings,
lifecycle: &ksp_worker_api::WorkerLifecycle,
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,
persistence: &mut PersistenceTasks,
port: &std::option::Option<PersistencePort>,
snapshots: &mut crate::RawTransactionIngestSnapshotPublisher,
) -> std::option::Option<ksp_core_lib::ErrorCode> {
let mut admission_open = true;
loop {
@@ -352,7 +434,7 @@ async fn supervise_until_stop(
}
joined = persistence.join_next(), if !persistence.is_empty() => {
if let std::option::Option::Some(value) = joined {
let fault = persistence_fault(value);
let fault = persistence_completion(value, lifecycle.state(), admission.queue_depth(), persistence.len(), snapshots);
if fault.is_some() {
source_stop_sender.send_replace(true);
return fault;
@@ -362,15 +444,25 @@ async fn supervise_until_stop(
received = admission.receive(settings.network()), if admission_open && persistence.len() < settings.persistence_concurrency() => {
match received {
std::result::Result::Ok(std::option::Option::Some(acquisition)) => {
if !spawn_persistence(persistence, port, acquisition) {
let spawned = spawn_persistence(persistence, port, acquisition);
let published = snapshots.record_admission_success(lifecycle.state(), admission.queue_depth(), persistence.len());
if !spawned {
source_stop_sender.send_replace(true);
return std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
}
if let std::result::Result::Err(error) = published {
source_stop_sender.send_replace(true);
return std::option::Option::Some(error.code());
}
},
std::result::Result::Ok(std::option::Option::None) => admission_open = false,
std::result::Result::Err(_) => {
let published = snapshots.record_admission_failure(lifecycle.state(), admission.queue_depth(), persistence.len());
source_stop_sender.send_replace(true);
return std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
return match published {
std::result::Result::Ok(()) => std::option::Option::Some(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID),
std::result::Result::Err(error) => std::option::Option::Some(error.code()),
};
},
}
}