v0.3.8-pre.007
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-app-store-desk/src/store_runtime.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
//! Composite-selected Store readiness, status and bounded shutdown lifecycle owned by Store Desk.
|
||||
|
||||
@@ -35,6 +35,27 @@ impl crate::StoreStartup {
|
||||
};
|
||||
}
|
||||
|
||||
/// Executes one bounded server-side RAW transaction inspection query.
|
||||
pub(crate) async fn query_transactions(
|
||||
&self,
|
||||
request: crate::StoreTransactionQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<crate::StoreTransactionQueryResponseDto> {
|
||||
let runtime = self.runtime.as_ref();
|
||||
return match runtime {
|
||||
std::option::Option::Some(value) => value.query_transactions(request).await,
|
||||
std::option::Option::None => std::result::Result::Err(store_runtime_unavailable()),
|
||||
};
|
||||
}
|
||||
|
||||
/// Loads one explicit bounded RAW transaction detail projection.
|
||||
pub(crate) async fn transaction_detail(&self, request: crate::StoreTransactionDetailRequestDto) -> ksp_core_lib::Result<crate::StoreTransactionDetailDto> {
|
||||
let runtime = self.runtime.as_ref();
|
||||
return match runtime {
|
||||
std::option::Option::Some(value) => value.transaction_detail(request).await,
|
||||
std::option::Option::None => std::result::Result::Err(store_runtime_unavailable()),
|
||||
};
|
||||
}
|
||||
|
||||
/// Closes the retained Store runtime if startup reached the physical Store-open phase.
|
||||
pub(crate) async fn close(&self) -> ksp_core_lib::Result<()> {
|
||||
let runtime = self.runtime.as_ref();
|
||||
@@ -95,6 +116,68 @@ impl crate::StoreRuntime {
|
||||
};
|
||||
}
|
||||
|
||||
/// Executes one validated random-access transaction inspection query while retaining the Store read guard.
|
||||
pub(crate) async fn query_transactions(
|
||||
&self,
|
||||
request: crate::StoreTransactionQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<crate::StoreTransactionQueryResponseDto> {
|
||||
let locked = self.store.read().await;
|
||||
let store = locked.as_ref();
|
||||
let store = match store {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return std::result::Result::Err(store_runtime_unavailable()),
|
||||
};
|
||||
let query = transaction_query_from_request(store, &request);
|
||||
let query = match query {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let page = ksp_store_lib::RawTransactionInspectionRead::inspect_raw_transactions(store, &query).await;
|
||||
let page = match page {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let records_total_decimal = page.total_items().to_string();
|
||||
let records_filtered_decimal = page.filtered_items().to_string();
|
||||
let items = page.into_items();
|
||||
let mut rows = std::vec::Vec::with_capacity(items.len());
|
||||
for item in items {
|
||||
rows.push(transaction_row_dto(&item));
|
||||
}
|
||||
return std::result::Result::Ok(crate::StoreTransactionQueryResponseDto { records_filtered_decimal, records_total_decimal, rows });
|
||||
}
|
||||
|
||||
/// Loads one transaction detail by application-owned row identity without exposing Store network selection to the frontend.
|
||||
pub(crate) async fn transaction_detail(&self, request: crate::StoreTransactionDetailRequestDto) -> ksp_core_lib::Result<crate::StoreTransactionDetailDto> {
|
||||
let locked = self.store.read().await;
|
||||
let store = locked.as_ref();
|
||||
let store = match store {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return std::result::Result::Err(store_runtime_unavailable()),
|
||||
};
|
||||
let reference = transaction_reference_from_request(store, request.signature.as_str());
|
||||
let reference = match reference {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let retention = ksp_store_lib::RawTransactionRetentionRead::get_raw_transaction_retention_state(store, &reference).await;
|
||||
let retention = match retention {
|
||||
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
||||
std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Err(store_transaction_not_found()),
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
return match retention {
|
||||
ksp_store_lib::RawRetentionState::Full | ksp_store_lib::RawRetentionState::Archived => {
|
||||
transaction_detail_with_payload(store, &reference, retention).await
|
||||
},
|
||||
ksp_store_lib::RawRetentionState::Purged => transaction_detail_from_tombstone(store, &reference).await,
|
||||
_ => std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_APP_STATE_INVALID,
|
||||
"Store Desk cannot project an unsupported transaction retention state",
|
||||
)),
|
||||
};
|
||||
}
|
||||
|
||||
/// Explicitly closes the Store exactly once through its backend-neutral facade after all read guards have drained.
|
||||
pub(crate) async fn close(&self) -> ksp_core_lib::Result<()> {
|
||||
let mut locked = self.store.write().await;
|
||||
@@ -127,6 +210,206 @@ impl crate::StoreRuntime {
|
||||
}
|
||||
}
|
||||
|
||||
const TRANSACTION_DETAIL_PREVIEW_BYTES: usize = 512;
|
||||
|
||||
fn transaction_query_from_request(
|
||||
store: &ksp_store_lib::Store,
|
||||
request: &crate::StoreTransactionQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<ksp_store_lib::RawTransactionInspectionQuery> {
|
||||
let limit = match request.limit {
|
||||
25 | 50 | 100 => ksp_store_lib::RawPageLimit::new(u64::from(request.limit)),
|
||||
_ => return std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
let limit = match limit {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
let offset = request.offset.parse::<u64>();
|
||||
let offset = match offset {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
let slot_min = parse_optional_decimal_u64(request.slot_min.as_deref());
|
||||
let slot_min = match slot_min {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let slot_max = parse_optional_decimal_u64(request.slot_max.as_deref());
|
||||
let slot_max = match slot_max {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let slots = ksp_store_lib::RawSlotRange::new(slot_min, slot_max);
|
||||
let slots = match slots {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
let direction = match request.direction.as_str() {
|
||||
"ascending" => ksp_store_lib::RawSortDirection::Ascending,
|
||||
"descending" => ksp_store_lib::RawSortDirection::Descending,
|
||||
_ => return std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
let network = store.runtime_snapshot().network().clone();
|
||||
let page = ksp_store_lib::RawInspectionPageRequest::new(offset, limit);
|
||||
return std::result::Result::Ok(ksp_store_lib::RawTransactionInspectionQuery::new(network, slots, direction, page));
|
||||
}
|
||||
|
||||
fn transaction_reference_from_request(store: &ksp_store_lib::Store, signature_text: &str) -> ksp_core_lib::Result<ksp_store_lib::RawTransactionReference> {
|
||||
let signature = decode_hex_64(signature_text);
|
||||
let signature = match signature {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let network = store.runtime_snapshot().network().clone();
|
||||
return std::result::Result::Ok(ksp_store_lib::RawTransactionReference::new(network, ksp_store_lib::RawTransactionSignature::new(signature)));
|
||||
}
|
||||
|
||||
async fn transaction_detail_with_payload(
|
||||
store: &ksp_store_lib::Store,
|
||||
reference: &ksp_store_lib::RawTransactionReference,
|
||||
retention: ksp_store_lib::RawRetentionState,
|
||||
) -> ksp_core_lib::Result<crate::StoreTransactionDetailDto> {
|
||||
let transaction = ksp_store_lib::RawTransactionRead::get_raw_transaction(store, reference).await;
|
||||
let transaction = match transaction {
|
||||
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
||||
std::result::Result::Ok(std::option::Option::None) => {
|
||||
return std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_APP_STATE_INVALID,
|
||||
"Store Desk transaction retention state and payload availability disagree",
|
||||
));
|
||||
},
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let payload = transaction.payload();
|
||||
let payload_bytes = payload.bytes();
|
||||
let preview_len = std::cmp::min(payload_bytes.len(), TRANSACTION_DETAIL_PREVIEW_BYTES);
|
||||
let preview = encode_hex(&payload_bytes[..preview_len]);
|
||||
return std::result::Result::Ok(crate::StoreTransactionDetailDto {
|
||||
block_time_unix_millis_decimal: transaction.block_time().map(|value| return value.unix_millis().to_string()),
|
||||
content_hash: encode_hex(payload.content_hash().as_bytes()),
|
||||
format_id: payload.format_id().as_str().to_owned(),
|
||||
format_version: payload.format_version(),
|
||||
payload_preview_hex: std::option::Option::Some(preview),
|
||||
payload_preview_truncated: payload_bytes.len() > preview_len,
|
||||
payload_size_decimal: std::option::Option::Some(payload.byte_len().to_string()),
|
||||
retention_state: retention_state_code(retention).to_owned(),
|
||||
signature: encode_hex(reference.signature().as_bytes()),
|
||||
slot_decimal: transaction.slot().to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
async fn transaction_detail_from_tombstone(
|
||||
store: &ksp_store_lib::Store,
|
||||
reference: &ksp_store_lib::RawTransactionReference,
|
||||
) -> ksp_core_lib::Result<crate::StoreTransactionDetailDto> {
|
||||
let tombstone = ksp_store_lib::RawTransactionRetentionRead::get_raw_transaction_tombstone(store, reference).await;
|
||||
let tombstone = match tombstone {
|
||||
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
||||
std::result::Result::Ok(std::option::Option::None) => {
|
||||
return std::result::Result::Err(ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_APP_STATE_INVALID,
|
||||
"Store Desk purged transaction is missing its durable tombstone",
|
||||
));
|
||||
},
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
return std::result::Result::Ok(crate::StoreTransactionDetailDto {
|
||||
block_time_unix_millis_decimal: std::option::Option::None,
|
||||
content_hash: encode_hex(tombstone.content_hash().as_bytes()),
|
||||
format_id: tombstone.format_id().as_str().to_owned(),
|
||||
format_version: tombstone.format_version(),
|
||||
payload_preview_hex: std::option::Option::None,
|
||||
payload_preview_truncated: false,
|
||||
payload_size_decimal: std::option::Option::None,
|
||||
retention_state: retention_state_code(ksp_store_lib::RawRetentionState::Purged).to_owned(),
|
||||
signature: encode_hex(reference.signature().as_bytes()),
|
||||
slot_decimal: tombstone.slot().to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
fn transaction_row_dto(summary: &ksp_store_lib::RawTransactionSummary) -> crate::StoreTransactionRowDto {
|
||||
return crate::StoreTransactionRowDto {
|
||||
block_time_unix_millis_decimal: summary.block_time().map(|value| return value.unix_millis().to_string()),
|
||||
content_hash: encode_hex(summary.content_hash().as_bytes()),
|
||||
format_id: summary.format_id().as_str().to_owned(),
|
||||
format_version: summary.format_version(),
|
||||
payload_size_decimal: summary.payload_size_bytes().map(|value| return value.to_string()),
|
||||
retention_state: retention_state_code(summary.retention_state()).to_owned(),
|
||||
signature: encode_hex(summary.reference().signature().as_bytes()),
|
||||
slot_decimal: summary.slot().to_string(),
|
||||
};
|
||||
}
|
||||
|
||||
fn retention_state_code(state: ksp_store_lib::RawRetentionState) -> &'static str {
|
||||
return match state {
|
||||
ksp_store_lib::RawRetentionState::Full => "full",
|
||||
ksp_store_lib::RawRetentionState::Compacted => "compacted",
|
||||
ksp_store_lib::RawRetentionState::Archived => "archived",
|
||||
ksp_store_lib::RawRetentionState::Purged => "purged",
|
||||
_ => "unknown",
|
||||
};
|
||||
}
|
||||
|
||||
fn parse_optional_decimal_u64(value: std::option::Option<&str>) -> ksp_core_lib::Result<std::option::Option<u64>> {
|
||||
return match value {
|
||||
std::option::Option::Some(text) => match text.parse::<u64>() {
|
||||
std::result::Result::Ok(parsed) => std::result::Result::Ok(std::option::Option::Some(parsed)),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(store_query_invalid()),
|
||||
},
|
||||
std::option::Option::None => std::result::Result::Ok(std::option::Option::None),
|
||||
};
|
||||
}
|
||||
|
||||
fn decode_hex_64(value: &str) -> ksp_core_lib::Result<[u8; 64]> {
|
||||
if value.len() != 128 || !value.is_ascii() {
|
||||
return std::result::Result::Err(store_query_invalid());
|
||||
}
|
||||
let source = value.as_bytes();
|
||||
let mut decoded = [0_u8; 64];
|
||||
for index in 0..64 {
|
||||
let high = decode_hex_nibble(source[index * 2]);
|
||||
let low = decode_hex_nibble(source[index * 2 + 1]);
|
||||
let (high, low) = match (high, low) {
|
||||
(std::option::Option::Some(high), std::option::Option::Some(low)) => (high, low),
|
||||
_ => return std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
decoded[index] = (high << 4) | low;
|
||||
}
|
||||
return std::result::Result::Ok(decoded);
|
||||
}
|
||||
|
||||
fn decode_hex_nibble(value: u8) -> std::option::Option<u8> {
|
||||
return match value {
|
||||
b'0'..=b'9' => std::option::Option::Some(value - b'0'),
|
||||
b'a'..=b'f' => std::option::Option::Some(value - b'a' + 10),
|
||||
b'A'..=b'F' => std::option::Option::Some(value - b'A' + 10),
|
||||
_ => std::option::Option::None,
|
||||
};
|
||||
}
|
||||
|
||||
fn encode_hex(bytes: &[u8]) -> String {
|
||||
const HEX: &[u8; 16] = b"0123456789abcdef";
|
||||
let mut encoded = String::with_capacity(bytes.len() * 2);
|
||||
for byte in bytes {
|
||||
let value = *byte;
|
||||
encoded.push(char::from(HEX[usize::from(value >> 4)]));
|
||||
encoded.push(char::from(HEX[usize::from(value & 0x0f)]));
|
||||
}
|
||||
return encoded;
|
||||
}
|
||||
|
||||
fn store_query_invalid() -> ksp_core_lib::Error {
|
||||
return ksp_core_lib::Error::new(crate::ERROR_CODE_STORE_QUERY_INVALID, "Store Desk transaction query is invalid");
|
||||
}
|
||||
|
||||
fn store_runtime_unavailable() -> ksp_core_lib::Error {
|
||||
return ksp_core_lib::Error::new(crate::ERROR_CODE_STORE_RUNTIME_UNAVAILABLE, "Store Desk Store runtime is unavailable");
|
||||
}
|
||||
|
||||
fn store_transaction_not_found() -> ksp_core_lib::Error {
|
||||
return ksp_core_lib::Error::new(crate::ERROR_CODE_STORE_TRANSACTION_NOT_FOUND, "Store Desk transaction was not found");
|
||||
}
|
||||
|
||||
/// Resolves the composite Store target, opens Store through the facade and captures one initial readiness probe.
|
||||
pub(crate) async fn initialize_store(management: &ksp_config_lib::ConfigManagement) -> crate::StoreStartup {
|
||||
let environment = ksp_config_lib::ConfigEnvironment::load();
|
||||
|
||||
Reference in New Issue
Block a user