// file: crates/ksp-app-raw-transaction-ingest-desk/src/route_runtime.rs // version: 8 //! Multi-route Store and independent Worker lifecycle owned by Raw Transaction Ingest Desk. const ROUTE_SHUTDOWN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15); const ROUTE_STOP_CLEANUP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15); 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); struct RouteRuntimeInner { last_monitoring: std::vec::Vec, last_terminals: std::vec::Vec, network: std::option::Option, next_sequence: u64, routes: std::vec::Vec, shutting_down: bool, store: std::option::Option>, store_closing_token: std::option::Option, store_opening_token: std::option::Option, } /// Shared multi-route runtime state. Each route owns an independent Worker while same-network routes share one Store. pub(crate) struct RouteRuntimeState { inner: std::sync::Mutex, } impl crate::RouteRuntimeState { /// Creates an idle multi-route runtime state without Store or Worker ownership. #[must_use] pub(crate) fn new() -> Self { return Self { inner: std::sync::Mutex::new(RouteRuntimeInner { last_monitoring: std::vec::Vec::new(), last_terminals: std::vec::Vec::new(), network: std::option::Option::None, next_sequence: 0, routes: std::vec::Vec::new(), shutting_down: false, store: std::option::Option::None, store_closing_token: std::option::Option::None, store_opening_token: std::option::Option::None, }), }; } /// Rejects new route Start work after application shutdown has closed admission. pub(crate) fn ensure_start_admission_open(&self) -> ksp_core_lib::Result<()> { 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()), }; if inner.shutting_down { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_SHUTTING_DOWN, "Raw Transaction Ingest Desk cannot start a route while application shutdown is in progress", )); } return std::result::Result::Ok(()); } fn reserve(&self, prepared: &crate::PreparedRouteStart) -> ksp_core_lib::Result { 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 inner.shutting_down { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_SHUTTING_DOWN, "Raw Transaction Ingest Desk cannot start a route while application shutdown is in progress", )); } if inner.store_closing_token.is_some() { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_ACTIVE, "Raw Transaction Ingest Desk Store is still closing after the previous route set", )); } if inner.routes.iter().any(|slot| return slot.matches_identity(prepared.profile_id.as_str(), prepared.route_id)) { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_ACTIVE, "Raw Transaction Ingest Desk already owns one Worker for the selected logical route", )); } if let std::option::Option::Some(network) = inner.network.as_deref() && network != prepared.network.as_str() { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_NETWORK_MISMATCH, "Raw Transaction Ingest Desk cannot share one Store across different logical networks", )); } let needs_store_open = inner.store.is_none(); if needs_store_open && inner.store_opening_token.is_some() { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_ACTIVE, "Raw Transaction Ingest Desk same-network Store is still opening", )); } let sequence = match inner.next_sequence.checked_add(1) { 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, }; if inner.network.is_none() { inner.last_monitoring.clear(); inner.last_terminals.clear(); inner.network = std::option::Option::Some(prepared.network.clone()); } inner.last_monitoring.retain(|status| return !identity.matches_monitoring(status)); inner.last_terminals.retain(|terminal| return !identity.matches_dto(terminal)); inner.next_sequence = sequence; if needs_store_open { inner.store_opening_token = std::option::Option::Some(sequence); } inner.routes.push(RouteRuntimeSlot::Starting { identity: identity.clone(), stop_requested: false, token: sequence }); return std::result::Result::Ok(RouteRuntimeReservation { identity, needs_store_open, token: sequence, worker_id }); } fn install_store(&self, token: u64, store: std::sync::Arc) -> ksp_core_lib::Result<()> { 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 inner.store_opening_token != std::option::Option::Some(token) || inner.store.is_some() { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_STATE_INVALID, "Raw Transaction Ingest Desk shared Store reservation changed before installation", )); } inner.store = std::option::Option::Some(store); inner.store_opening_token = std::option::Option::None; return std::result::Result::Ok(()); } fn shared_store(&self, expected_network: &str) -> ksp_core_lib::Result> { 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()), }; if inner.network.as_deref() != std::option::Option::Some(expected_network) { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_NETWORK_MISMATCH, "Raw Transaction Ingest Desk shared Store network no longer matches the selected route", )); } return match &inner.store { std::option::Option::Some(store) => std::result::Result::Ok(std::sync::Arc::clone(store)), std::option::Option::None => std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_STATE_INVALID, "Raw Transaction Ingest Desk shared Store is not installed", )), }; } fn activate( &self, reservation: &RouteRuntimeReservation, handle: ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHandle, ) -> ksp_core_lib::Result { 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()), }; let position = inner.routes.iter().position(|slot| return slot.starting_token() == std::option::Option::Some(reservation.token)); let position = match position { 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_STATE_INVALID, "Raw Transaction Ingest Desk route reservation changed before Worker activation", )); }, }; let stop_requested = match &inner.routes[position] { RouteRuntimeSlot::Starting { stop_requested, .. } => *stop_requested, RouteRuntimeSlot::Active { .. } => false, }; let should_stop = stop_requested || inner.shutting_down; if should_stop { let _accepted = handle.request_stop(); } let state = if should_stop { crate::RawIngestRouteState::Stopping } else { project_worker_state(handle.snapshot_source().current().worker_snapshot().state()) }; inner.routes[position] = RouteRuntimeSlot::Active { handle, identity: reservation.identity.clone(), token: reservation.token }; return std::result::Result::Ok(reservation.identity.dto(state)); } fn rollback(&self, token: u64) -> ksp_core_lib::Result>> { 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()), }; let cancelled_terminal = inner.routes.iter().find(|slot| return slot.token() == token).and_then(RouteRuntimeSlot::cancelled_start_terminal); inner.routes.retain(|slot| return slot.token() != token); if let std::option::Option::Some(terminal) = cancelled_terminal { inner.last_terminals.retain(|previous| return previous.profile_id != terminal.profile_id || previous.route_id != terminal.route_id); inner.last_terminals.push(terminal); } if inner.store_opening_token == std::option::Option::Some(token) { inner.store_opening_token = std::option::Option::None; } if inner.routes.is_empty() { let store = inner.store.take(); if store.is_some() { inner.store_closing_token = std::option::Option::Some(token); } else { inner.network = std::option::Option::None; } return std::result::Result::Ok(store); } return std::result::Result::Ok(std::option::Option::None); } fn finish( &self, token: u64, terminal: crate::RawIngestRouteRuntimeDto, monitoring: std::option::Option, ) -> ksp_core_lib::Result>> { 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()), }; let before = inner.routes.len(); inner.routes.retain(|slot| return slot.token() != token); if inner.routes.len() == before { return std::result::Result::Ok(std::option::Option::None); } inner.last_terminals.retain(|previous| return previous.profile_id != terminal.profile_id || previous.route_id != terminal.route_id); inner.last_terminals.push(terminal); if let std::option::Option::Some(status) = monitoring { inner.last_monitoring.retain(|previous| return previous.profile_id != status.profile_id || previous.route_id != status.route_id); inner.last_monitoring.push(status); } if inner.routes.is_empty() { let store = inner.store.take(); if store.is_some() { inner.store_closing_token = std::option::Option::Some(token); } else { inner.network = std::option::Option::None; } return std::result::Result::Ok(store); } return std::result::Result::Ok(std::option::Option::None); } fn complete_store_close(&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 inner.store_closing_token == std::option::Option::Some(token) { inner.store_closing_token = std::option::Option::None; if inner.routes.is_empty() && inner.store.is_none() && inner.store_opening_token.is_none() { inner.network = std::option::Option::None; } } } /// Returns complete latest-value monitoring projections for active and recent terminal routes in the current runtime session. pub(crate) fn monitoring_statuses(&self) -> ksp_core_lib::Result> { 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()), }; let mut values = inner.last_monitoring.clone(); for slot in &inner.routes { if let RouteRuntimeSlot::Active { handle, identity, .. } = slot { let snapshot = handle.snapshot_source().current(); let runtime = identity.dto(project_worker_state(snapshot.worker_snapshot().state())); let projected = crate::project_route_monitoring(&runtime, &snapshot); let projected = match projected { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; values.retain(|previous| return previous.profile_id != projected.profile_id || previous.route_id != projected.route_id); values.push(projected); } } return std::result::Result::Ok(values); } /// Closes route Start admission exactly once and requests cooperative Stop for all active or still-starting routes. pub(crate) fn begin_shutdown(&self) -> ksp_core_lib::Result { 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 inner.shutting_down { return std::result::Result::Ok(false); } inner.shutting_down = true; for slot in &mut inner.routes { match slot { RouteRuntimeSlot::Starting { stop_requested, .. } => *stop_requested = true, RouteRuntimeSlot::Active { handle, .. } => { let _accepted = handle.request_stop(); }, } } return std::result::Result::Ok(true); } /// Waits until all owned route Workers and the shared Store have completed bounded application-shutdown cleanup. pub(crate) async fn shutdown_and_wait(&self) -> ksp_core_lib::Result<()> { let begin = self.begin_shutdown(); if let std::result::Result::Err(error) = begin { return std::result::Result::Err(error); } let started = std::time::Instant::now(); loop { let complete = self.shutdown_complete(); let complete = match complete { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if complete { return std::result::Result::Ok(()); } if started.elapsed() >= ROUTE_SHUTDOWN_TIMEOUT { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_SHUTDOWN_FAILED, "Raw Transaction Ingest Desk application shutdown did not release all route Workers and shared Store state before the bounded deadline", )); } tokio::time::sleep(STORE_RECLAIM_POLL_INTERVAL).await; } } fn shutdown_complete(&self) -> ksp_core_lib::Result { 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( inner.routes.is_empty() && inner.store.is_none() && inner.store_opening_token.is_none() && inner.store_closing_token.is_none(), ); } /// Requests cooperative Stop for one exact logical route and waits for its Worker cleanup. pub(crate) async fn stop_and_wait(&self, request: &crate::RawIngestRouteStopRequestDto) -> ksp_core_lib::Result { let identity = request.validate_logical_identity(); if let std::result::Result::Err(error) = identity { return std::result::Result::Err(error); } let target = self.stop_target(request); let target = match target { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let (handle, identity, token) = match target { RouteStopTarget::Active { handle, identity, token } => (handle, identity, token), RouteStopTarget::Starting { identity, token } => return self.wait_route_cleanup(token, &identity).await, RouteStopTarget::Terminal(value) => return std::result::Result::Ok(value), RouteStopTarget::TerminalClosing => return self.wait_terminal_store_cleanup(request).await, }; let _first_stop_request = handle.request_stop(); if let std::result::Result::Err(error) = handle.wait_terminal().await { 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), ); } return self.wait_route_cleanup(token, &identity).await; } fn stop_target(&self, request: &crate::RawIngestRouteStopRequestDto) -> ksp_core_lib::Result { 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()), }; let position = inner.routes.iter().position(|slot| return slot.matches_identity(request.profile_id.as_str(), request.route_id)); if let std::option::Option::Some(position) = position { return match &mut inner.routes[position] { RouteRuntimeSlot::Active { handle, identity, token } => { std::result::Result::Ok(RouteStopTarget::Active { handle: handle.clone(), identity: identity.clone(), token: *token }) }, RouteRuntimeSlot::Starting { identity, stop_requested, token } => { *stop_requested = true; std::result::Result::Ok(RouteStopTarget::Starting { identity: identity.clone(), token: *token }) }, }; } let terminal = inner.last_terminals.iter().find(|terminal| return terminal.profile_id == request.profile_id && terminal.route_id == request.route_id); return match terminal { std::option::Option::Some(value) if inner.store_closing_token.is_none() => std::result::Result::Ok(RouteStopTarget::Terminal(value.clone())), std::option::Option::Some(_) => std::result::Result::Ok(RouteStopTarget::TerminalClosing), std::option::Option::None => std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_ROUTE_RUNTIME_NOT_ACTIVE, "Raw Transaction Ingest Desk has no active Worker for the selected logical route", )), }; } async fn wait_route_cleanup(&self, token: u64, identity: &RouteRuntimeIdentity) -> ksp_core_lib::Result { let started = std::time::Instant::now(); loop { let cleanup = self.route_cleanup_complete(token, identity); let cleanup = match cleanup { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if let std::option::Option::Some(value) = cleanup { return std::result::Result::Ok(value); } 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 selected route", )); } tokio::time::sleep(STORE_RECLAIM_POLL_INTERVAL).await; } } async fn wait_terminal_store_cleanup(&self, request: &crate::RawIngestRouteStopRequestDto) -> ksp_core_lib::Result { let started = std::time::Instant::now(); loop { let terminal = self.late_terminal_after_store_cleanup(request); let terminal = match terminal { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if let std::option::Option::Some(value) = terminal { return std::result::Result::Ok(value); } 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 final shared Store cleanup did not complete for the selected terminal route", )); } tokio::time::sleep(STORE_RECLAIM_POLL_INTERVAL).await; } } fn late_terminal_after_store_cleanup( &self, request: &crate::RawIngestRouteStopRequestDto, ) -> ksp_core_lib::Result> { 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()), }; let terminal = inner.last_terminals.iter().find(|terminal| return terminal.profile_id == request.profile_id && terminal.route_id == request.route_id); let terminal = match terminal { 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_NOT_ACTIVE, "Raw Transaction Ingest Desk terminal route is no longer retained", )); }, }; if inner.store_closing_token.is_some() { return std::result::Result::Ok(std::option::Option::None); } return std::result::Result::Ok(std::option::Option::Some(terminal.clone())); } fn route_cleanup_complete( &self, token: u64, identity: &RouteRuntimeIdentity, ) -> ksp_core_lib::Result> { 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()), }; if inner.routes.iter().any(|slot| return slot.token() == token) || inner.store_closing_token == std::option::Option::Some(token) { return std::result::Result::Ok(std::option::Option::None); } let terminal = inner.last_terminals.iter().find(|value| return identity.matches_dto(value)); return std::result::Result::Ok(terminal.cloned()); } } /// Launch ownership returned only after shared Store readiness, Worker creation and route-handle installation have succeeded. pub(crate) struct RouteRuntimeLaunch { acknowledgement: crate::RawIngestRouteRuntimeDto, handle: ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHandle, runtime_state: std::sync::Arc, 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(); } /// Returns the complete safe logical route identity used to project Worker monitoring snapshots. #[must_use] pub(crate) fn monitoring_identity(&self) -> crate::RawIngestRouteRuntimeDto { return self.acknowledgement.clone(); } /// Returns an independent latest-value source for frontend-safe monitoring relay. #[must_use] pub(crate) fn monitoring_source(&self) -> ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSnapshotSource { return self.handle.snapshot_source(); } /// Waits for this Worker terminal state, removes only its route and closes Store only after the final same-network route terminates. pub(crate) async fn monitor(self) { let terminal = self.handle.wait_terminal().await; let snapshot = self.handle.snapshot_source().current(); let worker = snapshot.worker_snapshot(); let snapshot_fault = worker.state().fault_code(); let snapshot_fault_domain = snapshot_fault.map_or("none", |code| return code.domain()); let snapshot_fault_code = snapshot_fault.map_or("none", |code| return code.code()); let terminal_state = match terminal { std::result::Result::Ok(state) => { ksp_logging_lib::info!( target: crate::TRACING_TARGET, domain = crate::TRACING_DOMAIN_ROUTE_RUNTIME, profile_id = self.acknowledgement.profile_id.as_str(), network = self.acknowledgement.network.as_str(), route_id = self.acknowledgement.route_id.as_str(), commitment = self.acknowledgement.commitment.as_str(), worker_state = state.code(), worker_health = worker.health().code(), worker_activity = worker.activity().code(), fault_domain = snapshot_fault_domain, fault_code = snapshot_fault_code, source_failure_total = snapshot.source_failure_total(), store_failure_total = snapshot.store_failure_total(), content_conflict_total = snapshot.content_conflict_total(), source_active = snapshot.source_active(), source_reconnecting = snapshot.source_reconnecting(), source_failed = snapshot.source_failed(), open_gap_count = snapshot.open_gap_count(), continuity_has_open_gaps = snapshot.continuity_has_open_gaps(), future_target_coverage = snapshot.future_target_coverage(), failed_source_losses_reconciled = snapshot.failed_source_losses_reconciled(), "Raw Transaction Ingest Desk route Worker reached terminal state" ); project_worker_state(state) }, std::result::Result::Err(error) => { ksp_logging_lib::warn!( target: crate::TRACING_TARGET, domain = crate::TRACING_DOMAIN_ROUTE_RUNTIME, profile_id = self.acknowledgement.profile_id.as_str(), network = self.acknowledgement.network.as_str(), route_id = self.acknowledgement.route_id.as_str(), commitment = self.acknowledgement.commitment.as_str(), error_domain = error.code().domain(), error_code = error.code().code(), worker_health = worker.health().code(), snapshot_fault_domain = snapshot_fault_domain, snapshot_fault_code = snapshot_fault_code, source_failure_total = snapshot.source_failure_total(), store_failure_total = snapshot.store_failure_total(), content_conflict_total = snapshot.content_conflict_total(), "Raw Transaction Ingest Desk terminal Worker wait failed" ); crate::RawIngestRouteState::Faulted }, }; let mut terminal = self.acknowledgement; terminal.state = terminal_state; let monitoring = crate::project_route_monitoring(&terminal, &snapshot); let monitoring = match monitoring { std::result::Result::Ok(value) => std::option::Option::Some(value), 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 could not retain terminal route monitoring projection" ); std::option::Option::None }, }; let store = self.runtime_state.finish(self.token, terminal, monitoring); let store = match store { std::result::Result::Ok(value) => value, 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 route cleanup state update failed" ); return; }, }; if let std::option::Option::Some(store) = store { let close = close_store_arc(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 shared Store shutdown failed after final Worker" ); } self.runtime_state.complete_store_close(self.token); } } } /// Starts one exact route Worker and reuses the same ready Store for any already-active same-network route. pub(crate) async fn start_route_runtime( runtime_state: std::sync::Arc, prepared: crate::PreparedRouteStart, ) -> ksp_core_lib::Result { let reservation = match runtime_state.reserve(&prepared) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let store = if reservation.needs_store_open { 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) => { let _rollback = 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 shared Store runtime") .with_source(error), ); }, }; let health = store.health().await; if health.state() != ksp_store_lib::StoreHealthState::Ready { let close = store.close().await; let _rollback = 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 shared 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 shared Store readiness probe did not prove Ready", )); } let store = std::sync::Arc::new(store); if let std::result::Result::Err(error) = runtime_state.install_store(reservation.token, std::sync::Arc::clone(&store)) { let close = close_store_arc(store).await; let _rollback = 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); } store } else { std::mem::drop(prepared.store_settings); match runtime_state.shared_store(prepared.network.as_str()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => { let _rollback = runtime_state.rollback(reservation.token); return std::result::Result::Err(error); }, } }; let network = store.runtime_snapshot().network().clone(); let settings = ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings::with_defaults(network, reservation.worker_id.clone()); 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 cleanup_store = runtime_state.rollback(reservation.token); let cleanup_store = match cleanup_store { std::result::Result::Ok(value) => value, std::result::Result::Err(rollback_error) => return std::result::Result::Err(rollback_error), }; std::mem::drop(store); if let std::option::Option::Some(value) = cleanup_store { let close = close_store_arc(value).await; runtime_state.complete_store_close(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 route 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 cleanup_store = runtime_state.rollback(reservation.token); let cleanup_store = match cleanup_store { std::result::Result::Ok(value) => value, std::result::Result::Err(rollback_error) => return std::result::Result::Err(rollback_error), }; std::mem::drop(store); if let std::option::Option::Some(value) = cleanup_store { let close = close_store_arc(value).await; runtime_state.complete_store_close(reservation.token); if let std::result::Result::Err(close_error) = close { return std::result::Result::Err(close_error); } } return std::result::Result::Err(error); }, }; std::mem::drop(store); 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 independent route Worker handle installed on shared Store" ); return std::result::Result::Ok(crate::RouteRuntimeLaunch { acknowledgement, handle, runtime_state, 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, }; } fn matches_dto(&self, value: &crate::RawIngestRouteRuntimeDto) -> bool { return self.profile_id == value.profile_id && self.route_id == value.route_id; } fn matches_monitoring(&self, value: &crate::RawIngestRouteMonitoringDto) -> bool { return self.profile_id == value.profile_id && self.route_id == value.route_id; } } enum RouteStopTarget { Active { handle: ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHandle, identity: RouteRuntimeIdentity, token: u64, }, Starting { identity: RouteRuntimeIdentity, token: u64, }, Terminal(crate::RawIngestRouteRuntimeDto), TerminalClosing, } struct RouteRuntimeReservation { identity: RouteRuntimeIdentity, needs_store_open: bool, token: u64, worker_id: ksp_worker_api::WorkerId, } enum RouteRuntimeSlot { Starting { identity: RouteRuntimeIdentity, stop_requested: bool, token: u64, }, Active { handle: ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHandle, identity: RouteRuntimeIdentity, token: u64, }, } impl RouteRuntimeSlot { fn matches_identity(&self, profile_id: &str, route_id: crate::RawIngestRouteId) -> bool { let identity = match self { Self::Starting { identity, .. } | Self::Active { identity, .. } => identity, }; return identity.profile_id == profile_id && identity.route_id == route_id; } fn starting_token(&self) -> std::option::Option { return match self { Self::Starting { token, .. } => std::option::Option::Some(*token), Self::Active { .. } => std::option::Option::None, }; } fn cancelled_start_terminal(&self) -> std::option::Option { return match self { Self::Starting { identity, stop_requested: true, .. } => std::option::Option::Some(identity.dto(crate::RawIngestRouteState::Stopped)), Self::Starting { stop_requested: false, .. } | Self::Active { .. } => std::option::Option::None, }; } fn token(&self) -> u64 { return match self { Self::Starting { token, .. } | Self::Active { token, .. } => *token, }; } } async fn close_store_arc(store: std::sync::Arc) -> 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 shared 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 shared Store remained referenced after final 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 multi-route runtime state lock is poisoned"); } #[cfg(test)] #[path = "../unit_tests/route_runtime.rs"] mod tests;