v0.3.15-pre.008
This commit is contained in:
396
crates/ksp-app-raw-transaction-ingest-desk/src/route_runtime.rs
Normal file
396
crates/ksp-app-raw-transaction-ingest-desk/src/route_runtime.rs
Normal file
@@ -0,0 +1,396 @@
|
||||
// file: crates/ksp-app-raw-transaction-ingest-desk/src/route_runtime.rs
|
||||
// version: 1
|
||||
|
||||
//! Mono-route Store and Worker lifecycle owned by Raw Transaction Ingest Desk.
|
||||
|
||||
const ROUTE_STOP_CLEANUP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
|
||||
const STORE_RECLAIM_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(5);
|
||||
const STORE_RECLAIM_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(1);
|
||||
|
||||
/// Shared mono-route runtime state. `pre.008` admits at most one active or starting Worker.
|
||||
pub(crate) struct RouteRuntimeState {
|
||||
inner: std::sync::Mutex<RouteRuntimeInner>,
|
||||
}
|
||||
|
||||
impl crate::RouteRuntimeState {
|
||||
/// Creates an idle mono-route runtime state.
|
||||
#[must_use]
|
||||
pub(crate) fn new() -> Self {
|
||||
return Self { inner: std::sync::Mutex::new(RouteRuntimeInner { next_sequence: 0, slot: RouteRuntimeSlot::Idle }) };
|
||||
}
|
||||
|
||||
fn reserve(&self, prepared: &crate::PreparedRouteStart) -> ksp_core_lib::Result<RouteRuntimeReservation> {
|
||||
let inner = self.inner.lock();
|
||||
let mut inner = match inner {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(runtime_lock_error()),
|
||||
};
|
||||
if !matches!(inner.slot, RouteRuntimeSlot::Idle) {
|
||||
return std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_ACTIVE,
|
||||
"Raw Transaction Ingest Desk already owns one mono-route Worker runtime",
|
||||
));
|
||||
}
|
||||
let sequence = inner.next_sequence.checked_add(1);
|
||||
let sequence = match sequence {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => {
|
||||
return std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_START_FAILED,
|
||||
"Raw Transaction Ingest Desk Worker sequence is exhausted",
|
||||
));
|
||||
},
|
||||
};
|
||||
let worker_id = ksp_worker_api::WorkerId::new(format!("raw-ingest-desk-{sequence}"));
|
||||
let worker_id = match worker_id {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => {
|
||||
return std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(crate::ERROR_CODE_ROUTE_RUNTIME_START_FAILED, "Cannot create Raw Transaction Ingest Desk Worker identity")
|
||||
.with_source(error),
|
||||
);
|
||||
},
|
||||
};
|
||||
let identity = RouteRuntimeIdentity {
|
||||
commitment: prepared.commitment,
|
||||
inventory_generation: prepared.inventory_generation,
|
||||
network: prepared.network.clone(),
|
||||
profile_id: prepared.profile_id.clone(),
|
||||
route_id: prepared.route_id,
|
||||
};
|
||||
inner.next_sequence = sequence;
|
||||
inner.slot = RouteRuntimeSlot::Starting { token: sequence };
|
||||
return std::result::Result::Ok(RouteRuntimeReservation { identity, token: sequence, worker_id });
|
||||
}
|
||||
|
||||
fn activate(
|
||||
&self,
|
||||
reservation: &RouteRuntimeReservation,
|
||||
handle: ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHandle,
|
||||
) -> ksp_core_lib::Result<crate::RawIngestRouteRuntimeDto> {
|
||||
let inner = self.inner.lock();
|
||||
let mut inner = match inner {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(runtime_lock_error()),
|
||||
};
|
||||
match &inner.slot {
|
||||
RouteRuntimeSlot::Starting { token } if *token == reservation.token => {},
|
||||
RouteRuntimeSlot::Idle | RouteRuntimeSlot::Starting { .. } | RouteRuntimeSlot::Active { .. } => {
|
||||
return std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_STATE_INVALID,
|
||||
"Raw Transaction Ingest Desk mono-route runtime reservation changed before Worker activation",
|
||||
));
|
||||
},
|
||||
}
|
||||
let state = project_worker_state(handle.snapshot_source().current().worker_snapshot().state());
|
||||
inner.slot = RouteRuntimeSlot::Active { handle, identity: reservation.identity.clone(), token: reservation.token };
|
||||
return std::result::Result::Ok(reservation.identity.dto(state));
|
||||
}
|
||||
|
||||
fn rollback(&self, token: u64) {
|
||||
let inner = self.inner.lock();
|
||||
let mut inner = match inner {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return,
|
||||
};
|
||||
if matches!(&inner.slot, RouteRuntimeSlot::Starting { token: current, .. } if *current == token) {
|
||||
inner.slot = RouteRuntimeSlot::Idle;
|
||||
}
|
||||
}
|
||||
|
||||
fn finish(&self, token: u64) {
|
||||
let inner = self.inner.lock();
|
||||
let mut inner = match inner {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return,
|
||||
};
|
||||
if matches!(&inner.slot, RouteRuntimeSlot::Active { token: current, .. } if *current == token) {
|
||||
inner.slot = RouteRuntimeSlot::Idle;
|
||||
}
|
||||
}
|
||||
|
||||
/// Requests cooperative Stop, waits for terminal Worker state and returns after the terminal monitor has released the mono-route slot.
|
||||
pub(crate) async fn stop_and_wait(&self) -> ksp_core_lib::Result<crate::RawIngestRouteRuntimeDto> {
|
||||
let active = {
|
||||
let inner = self.inner.lock();
|
||||
let inner = match inner {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(runtime_lock_error()),
|
||||
};
|
||||
match &inner.slot {
|
||||
RouteRuntimeSlot::Active { handle, identity, .. } => std::result::Result::Ok((handle.clone(), identity.clone())),
|
||||
RouteRuntimeSlot::Idle | RouteRuntimeSlot::Starting { .. } => std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_NOT_ACTIVE,
|
||||
"Raw Transaction Ingest Desk has no active mono-route Worker to stop",
|
||||
)),
|
||||
}
|
||||
};
|
||||
let (handle, identity) = match active {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let _first_stop_request = handle.request_stop();
|
||||
let terminal = handle.wait_terminal().await;
|
||||
let terminal = match terminal {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => {
|
||||
return std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(crate::ERROR_CODE_ROUTE_RUNTIME_STATE_INVALID, "Cannot observe terminal Raw Transaction Ingest Desk Worker state")
|
||||
.with_source(error),
|
||||
);
|
||||
},
|
||||
};
|
||||
let started = std::time::Instant::now();
|
||||
loop {
|
||||
if self.is_idle()? {
|
||||
return std::result::Result::Ok(identity.dto(project_worker_state(terminal)));
|
||||
}
|
||||
if started.elapsed() >= ROUTE_STOP_CLEANUP_TIMEOUT {
|
||||
return std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_STORE_SHUTDOWN_FAILED,
|
||||
"Raw Transaction Ingest Desk terminal cleanup did not release the mono-route slot",
|
||||
));
|
||||
}
|
||||
tokio::time::sleep(STORE_RECLAIM_POLL_INTERVAL).await;
|
||||
}
|
||||
}
|
||||
|
||||
fn is_idle(&self) -> ksp_core_lib::Result<bool> {
|
||||
let inner = self.inner.lock();
|
||||
let inner = match inner {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(runtime_lock_error()),
|
||||
};
|
||||
return std::result::Result::Ok(matches!(inner.slot, RouteRuntimeSlot::Idle));
|
||||
}
|
||||
}
|
||||
|
||||
/// Launch ownership returned only after Store readiness, Worker creation and handle installation have succeeded.
|
||||
pub(crate) struct RouteRuntimeLaunch {
|
||||
acknowledgement: crate::RawIngestRouteRuntimeDto,
|
||||
handle: ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHandle,
|
||||
runtime_state: std::sync::Arc<crate::RouteRuntimeState>,
|
||||
store: std::sync::Arc<ksp_store_lib::Store>,
|
||||
token: u64,
|
||||
}
|
||||
|
||||
impl crate::RouteRuntimeLaunch {
|
||||
/// Returns the safe Start acknowledgement captured after the Worker handle was installed.
|
||||
#[must_use]
|
||||
pub(crate) fn acknowledgement(&self) -> crate::RawIngestRouteRuntimeDto {
|
||||
return self.acknowledgement.clone();
|
||||
}
|
||||
|
||||
/// Waits for terminal Worker state, explicitly closes Store after the Worker releases its `Arc`, then frees the mono-route slot.
|
||||
pub(crate) async fn monitor(self) {
|
||||
let terminal = self.handle.wait_terminal().await;
|
||||
match terminal {
|
||||
std::result::Result::Ok(state) => {
|
||||
ksp_logging_lib::info!(
|
||||
target: crate::TRACING_TARGET,
|
||||
domain = crate::TRACING_DOMAIN_ROUTE_RUNTIME,
|
||||
worker_state = state.code(),
|
||||
"Raw Transaction Ingest Desk mono-route Worker reached terminal state"
|
||||
);
|
||||
},
|
||||
std::result::Result::Err(error) => {
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
domain = crate::TRACING_DOMAIN_ROUTE_RUNTIME,
|
||||
error_domain = error.code().domain(),
|
||||
error_code = error.code().code(),
|
||||
"Raw Transaction Ingest Desk terminal Worker wait failed"
|
||||
);
|
||||
},
|
||||
}
|
||||
let close = close_store_arc(self.store).await;
|
||||
if let std::result::Result::Err(error) = close {
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
domain = crate::TRACING_DOMAIN_ROUTE_RUNTIME,
|
||||
error_domain = error.code().domain(),
|
||||
error_code = error.code().code(),
|
||||
"Raw Transaction Ingest Desk Store shutdown failed after terminal Worker"
|
||||
);
|
||||
}
|
||||
self.runtime_state.finish(self.token);
|
||||
}
|
||||
}
|
||||
|
||||
/// Opens one Store, starts one exact Worker route and atomically installs its handle in the mono-route application slot.
|
||||
pub(crate) async fn start_route_runtime(
|
||||
runtime_state: std::sync::Arc<crate::RouteRuntimeState>,
|
||||
prepared: crate::PreparedRouteStart,
|
||||
) -> ksp_core_lib::Result<crate::RouteRuntimeLaunch> {
|
||||
let reservation = runtime_state.reserve(&prepared);
|
||||
let reservation = match reservation {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let store = ksp_store_lib::Store::open(prepared.store_settings).await;
|
||||
let store = match store {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => {
|
||||
runtime_state.rollback(reservation.token);
|
||||
return std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(crate::ERROR_CODE_ROUTE_RUNTIME_START_FAILED, "Cannot open Raw Transaction Ingest Desk Store runtime")
|
||||
.with_source(error),
|
||||
);
|
||||
},
|
||||
};
|
||||
let health = store.health().await;
|
||||
if health.state() != ksp_store_lib::StoreHealthState::Ready {
|
||||
let close = store.close().await;
|
||||
runtime_state.rollback(reservation.token);
|
||||
if let std::result::Result::Err(error) = close {
|
||||
return std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_STORE_SHUTDOWN_FAILED,
|
||||
"Cannot close non-ready Raw Transaction Ingest Desk Store runtime",
|
||||
)
|
||||
.with_source(error),
|
||||
);
|
||||
}
|
||||
return std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_STORE_NOT_READY,
|
||||
"Raw Transaction Ingest Desk Store readiness probe did not prove Ready",
|
||||
));
|
||||
}
|
||||
let network = store.runtime_snapshot().network().clone();
|
||||
let settings = ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings::with_defaults(network, reservation.worker_id.clone());
|
||||
let store = std::sync::Arc::new(store);
|
||||
let handle = ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestWorker::start_with_runtime_resources(
|
||||
settings,
|
||||
std::sync::Arc::clone(&store),
|
||||
prepared.runtime_resources,
|
||||
);
|
||||
let handle = match handle {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => {
|
||||
let close = close_store_arc(store).await;
|
||||
runtime_state.rollback(reservation.token);
|
||||
if let std::result::Result::Err(close_error) = close {
|
||||
return std::result::Result::Err(close_error);
|
||||
}
|
||||
return std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(crate::ERROR_CODE_ROUTE_RUNTIME_START_FAILED, "Cannot start Raw Transaction Ingest Desk Worker runtime")
|
||||
.with_source(error),
|
||||
);
|
||||
},
|
||||
};
|
||||
let acknowledgement = runtime_state.activate(&reservation, handle.clone());
|
||||
let acknowledgement = match acknowledgement {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => {
|
||||
let _stop_requested = handle.request_stop();
|
||||
let _terminal = handle.wait_terminal().await;
|
||||
let close = close_store_arc(store).await;
|
||||
runtime_state.rollback(reservation.token);
|
||||
if let std::result::Result::Err(close_error) = close {
|
||||
return std::result::Result::Err(close_error);
|
||||
}
|
||||
return std::result::Result::Err(error);
|
||||
},
|
||||
};
|
||||
ksp_logging_lib::info!(
|
||||
target: crate::TRACING_TARGET,
|
||||
domain = crate::TRACING_DOMAIN_ROUTE_RUNTIME,
|
||||
profile_id = reservation.identity.profile_id.as_str(),
|
||||
network = reservation.identity.network.as_str(),
|
||||
route_id = reservation.identity.route_id.as_str(),
|
||||
commitment = reservation.identity.commitment.as_str(),
|
||||
"Raw Transaction Ingest Desk mono-route Worker handle installed"
|
||||
);
|
||||
return std::result::Result::Ok(crate::RouteRuntimeLaunch { acknowledgement, handle, runtime_state, store, token: reservation.token });
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct RouteRuntimeIdentity {
|
||||
commitment: crate::RawIngestCommitment,
|
||||
inventory_generation: u32,
|
||||
network: String,
|
||||
profile_id: String,
|
||||
route_id: crate::RawIngestRouteId,
|
||||
}
|
||||
|
||||
impl RouteRuntimeIdentity {
|
||||
fn dto(&self, state: crate::RawIngestRouteState) -> crate::RawIngestRouteRuntimeDto {
|
||||
return crate::RawIngestRouteRuntimeDto {
|
||||
commitment: self.commitment,
|
||||
inventory_generation: self.inventory_generation,
|
||||
network: self.network.clone(),
|
||||
profile_id: self.profile_id.clone(),
|
||||
route_id: self.route_id,
|
||||
state,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
struct RouteRuntimeReservation {
|
||||
identity: RouteRuntimeIdentity,
|
||||
token: u64,
|
||||
worker_id: ksp_worker_api::WorkerId,
|
||||
}
|
||||
|
||||
enum RouteRuntimeSlot {
|
||||
Idle,
|
||||
Starting {
|
||||
token: u64,
|
||||
},
|
||||
Active {
|
||||
handle: ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHandle,
|
||||
identity: RouteRuntimeIdentity,
|
||||
token: u64,
|
||||
},
|
||||
}
|
||||
|
||||
async fn close_store_arc(store: std::sync::Arc<ksp_store_lib::Store>) -> ksp_core_lib::Result<()> {
|
||||
let started = std::time::Instant::now();
|
||||
let mut store = store;
|
||||
loop {
|
||||
match std::sync::Arc::try_unwrap(store) {
|
||||
std::result::Result::Ok(value) => {
|
||||
let close = value.close().await;
|
||||
return match close {
|
||||
std::result::Result::Ok(()) => std::result::Result::Ok(()),
|
||||
std::result::Result::Err(error) => std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_STORE_SHUTDOWN_FAILED,
|
||||
"Cannot close Raw Transaction Ingest Desk Store runtime",
|
||||
)
|
||||
.with_source(error),
|
||||
),
|
||||
};
|
||||
},
|
||||
std::result::Result::Err(shared) => {
|
||||
store = shared;
|
||||
},
|
||||
}
|
||||
if started.elapsed() >= STORE_RECLAIM_TIMEOUT {
|
||||
return std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_ROUTE_RUNTIME_STORE_SHUTDOWN_FAILED,
|
||||
"Raw Transaction Ingest Desk Store remained shared after terminal Worker",
|
||||
));
|
||||
}
|
||||
tokio::time::sleep(STORE_RECLAIM_POLL_INTERVAL).await;
|
||||
}
|
||||
}
|
||||
|
||||
fn project_worker_state(state: ksp_worker_api::WorkerState) -> crate::RawIngestRouteState {
|
||||
return match state {
|
||||
ksp_worker_api::WorkerState::Created | ksp_worker_api::WorkerState::Starting => crate::RawIngestRouteState::Starting,
|
||||
ksp_worker_api::WorkerState::Running => crate::RawIngestRouteState::Running,
|
||||
ksp_worker_api::WorkerState::Stopping => crate::RawIngestRouteState::Stopping,
|
||||
ksp_worker_api::WorkerState::Stopped => crate::RawIngestRouteState::Stopped,
|
||||
ksp_worker_api::WorkerState::Faulted(_) => crate::RawIngestRouteState::Faulted,
|
||||
_ => crate::RawIngestRouteState::Faulted,
|
||||
};
|
||||
}
|
||||
|
||||
fn runtime_lock_error() -> ksp_core_lib::Error {
|
||||
return ksp_core_lib::Error::new(crate::ERROR_CODE_APP_STATE_LOCK_FAILED, "Raw Transaction Ingest Desk mono-route runtime state lock is poisoned");
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "../unit_tests/route_runtime.rs"]
|
||||
mod tests;
|
||||
Reference in New Issue
Block a user