v0.3.8-pre.010
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-app-store-desk/src/store_runtime.rs
|
||||
// version: 4
|
||||
// version: 5
|
||||
|
||||
//! Composite-selected Store readiness, status and bounded shutdown lifecycle owned by Store Desk.
|
||||
|
||||
@@ -53,6 +53,18 @@ impl crate::StoreStartup {
|
||||
};
|
||||
}
|
||||
|
||||
/// Executes one bounded metadata-only observation inspection for an explicit RAW account state.
|
||||
pub(crate) async fn query_account_observations(
|
||||
&self,
|
||||
request: crate::StoreAccountObservationQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<crate::StoreAccountObservationQueryResponseDto> {
|
||||
let runtime = self.runtime.as_ref();
|
||||
return match runtime {
|
||||
std::option::Option::Some(value) => value.query_account_observations(request).await,
|
||||
std::option::Option::None => std::result::Result::Err(store_runtime_unavailable()),
|
||||
};
|
||||
}
|
||||
|
||||
/// Executes one bounded server-side RAW transaction inspection query.
|
||||
pub(crate) async fn query_transactions(
|
||||
&self,
|
||||
@@ -74,6 +86,18 @@ impl crate::StoreStartup {
|
||||
};
|
||||
}
|
||||
|
||||
/// Executes one bounded metadata-only observation inspection for an explicit RAW transaction.
|
||||
pub(crate) async fn query_transaction_observations(
|
||||
&self,
|
||||
request: crate::StoreTransactionObservationQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<crate::StoreTransactionObservationQueryResponseDto> {
|
||||
let runtime = self.runtime.as_ref();
|
||||
return match runtime {
|
||||
std::option::Option::Some(value) => value.query_transaction_observations(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();
|
||||
@@ -184,6 +208,37 @@ impl crate::StoreRuntime {
|
||||
return std::result::Result::Ok(account_detail_dto(&state));
|
||||
}
|
||||
|
||||
/// Executes one validated metadata-only observation query for one exact account-state identity.
|
||||
pub(crate) async fn query_account_observations(
|
||||
&self,
|
||||
request: crate::StoreAccountObservationQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<crate::StoreAccountObservationQueryResponseDto> {
|
||||
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 = account_observation_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::RawAccountObservationInspectionRead::inspect_raw_account_observations(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(account_observation_row_dto(&item));
|
||||
}
|
||||
return std::result::Result::Ok(crate::StoreAccountObservationQueryResponseDto { records_filtered_decimal, records_total_decimal, rows });
|
||||
}
|
||||
|
||||
/// Executes one validated random-access transaction inspection query while retaining the Store read guard.
|
||||
pub(crate) async fn query_transactions(
|
||||
&self,
|
||||
@@ -246,6 +301,37 @@ impl crate::StoreRuntime {
|
||||
};
|
||||
}
|
||||
|
||||
/// Executes one validated metadata-only observation query for one exact transaction identity.
|
||||
pub(crate) async fn query_transaction_observations(
|
||||
&self,
|
||||
request: crate::StoreTransactionObservationQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<crate::StoreTransactionObservationQueryResponseDto> {
|
||||
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_observation_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::RawTransactionObservationInspectionRead::inspect_raw_transaction_observations(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_observation_row_dto(&item));
|
||||
}
|
||||
return std::result::Result::Ok(crate::StoreTransactionObservationQueryResponseDto { records_filtered_decimal, records_total_decimal, rows });
|
||||
}
|
||||
|
||||
/// 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;
|
||||
@@ -281,23 +367,88 @@ impl crate::StoreRuntime {
|
||||
const ACCOUNT_DETAIL_PREVIEW_BYTES: usize = 512;
|
||||
const TRANSACTION_DETAIL_PREVIEW_BYTES: usize = 512;
|
||||
|
||||
fn account_query_from_request(
|
||||
store: &ksp_store_lib::Store,
|
||||
request: &crate::StoreAccountQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<ksp_store_lib::RawAccountStateInspectionQuery> {
|
||||
let limit = match request.limit {
|
||||
25 | 50 | 100 => ksp_store_lib::RawPageLimit::new(u64::from(request.limit)),
|
||||
fn inspection_page_from_request(offset_text: &str, limit_value: u32) -> ksp_core_lib::Result<ksp_store_lib::RawInspectionPageRequest> {
|
||||
let limit = match limit_value {
|
||||
25 | 50 | 100 => ksp_store_lib::RawPageLimit::new(u64::from(limit_value)),
|
||||
_ => 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 = offset_text.parse::<u64>();
|
||||
let offset = match offset {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
return std::result::Result::Ok(ksp_store_lib::RawInspectionPageRequest::new(offset, limit));
|
||||
}
|
||||
|
||||
fn account_observation_query_from_request(
|
||||
store: &ksp_store_lib::Store,
|
||||
request: &crate::StoreAccountObservationQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<ksp_store_lib::RawAccountObservationInspectionQuery> {
|
||||
let page = inspection_page_from_request(request.offset.as_str(), request.limit);
|
||||
let page = match page {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let detail_request =
|
||||
crate::StoreAccountDetailRequestDto { pubkey: request.pubkey.clone(), slot: request.slot.clone(), state_hash: request.state_hash.clone() };
|
||||
let reference = account_reference_from_request(store, &detail_request);
|
||||
let reference = match reference {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let network = store.runtime_snapshot().network().clone();
|
||||
let query = ksp_store_lib::RawAccountObservationInspectionQuery::try_new(
|
||||
network,
|
||||
std::option::Option::Some(reference),
|
||||
ksp_store_lib::RawSortDirection::Descending,
|
||||
page,
|
||||
);
|
||||
return match query {
|
||||
std::result::Result::Ok(value) => std::result::Result::Ok(value),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
}
|
||||
|
||||
fn transaction_observation_query_from_request(
|
||||
store: &ksp_store_lib::Store,
|
||||
request: &crate::StoreTransactionObservationQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<ksp_store_lib::RawTransactionObservationInspectionQuery> {
|
||||
let page = inspection_page_from_request(request.offset.as_str(), request.limit);
|
||||
let page = match page {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
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 network = store.runtime_snapshot().network().clone();
|
||||
let query = ksp_store_lib::RawTransactionObservationInspectionQuery::try_new(
|
||||
network,
|
||||
std::option::Option::Some(reference),
|
||||
ksp_store_lib::RawSortDirection::Descending,
|
||||
page,
|
||||
);
|
||||
return match query {
|
||||
std::result::Result::Ok(value) => std::result::Result::Ok(value),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(store_query_invalid()),
|
||||
};
|
||||
}
|
||||
|
||||
fn account_query_from_request(
|
||||
store: &ksp_store_lib::Store,
|
||||
request: &crate::StoreAccountQueryRequestDto,
|
||||
) -> ksp_core_lib::Result<ksp_store_lib::RawAccountStateInspectionQuery> {
|
||||
let page = inspection_page_from_request(request.offset.as_str(), request.limit);
|
||||
let page = match page {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let pubkey = parse_optional_pubkey(request.pubkey.as_deref());
|
||||
let pubkey = match pubkey {
|
||||
std::result::Result::Ok(value) => value,
|
||||
@@ -324,7 +475,6 @@ fn account_query_from_request(
|
||||
_ => 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::RawAccountStateInspectionQuery::new(network, pubkey, slots, direction, page));
|
||||
}
|
||||
|
||||
@@ -383,22 +533,59 @@ fn account_row_dto(summary: &ksp_store_lib::RawAccountStateSummary) -> crate::St
|
||||
};
|
||||
}
|
||||
|
||||
fn observation_provenance_dto(provenance: &ksp_store_lib::RawAcquisitionProvenance) -> crate::StoreObservationProvenanceDto {
|
||||
return crate::StoreObservationProvenanceDto {
|
||||
acquisition_method: provenance.acquisition_method().as_str().to_owned(),
|
||||
capture_session_id: provenance.capture_session_id().map(|value| return value.as_str().to_owned()),
|
||||
commitment: provenance.commitment().map(|value| return value.as_str().to_owned()),
|
||||
endpoint_id: provenance.endpoint_id().map(|value| return value.as_str().to_owned()),
|
||||
filter_id: provenance.filter_id().map(|value| return value.as_str().to_owned()),
|
||||
observed_at_unix_millis_decimal: provenance.observed_at().map(|value| return value.unix_millis().to_string()),
|
||||
origin: acquisition_origin_code(provenance.origin()).to_owned(),
|
||||
protocol: provenance.protocol().as_str().to_owned(),
|
||||
provider: provenance.provider().as_str().to_owned(),
|
||||
received_at_unix_millis_decimal: provenance.received_at().unix_millis().to_string(),
|
||||
source_payload_hash: provenance.source_payload_hash().map(|value| return encode_hex(value.as_bytes())),
|
||||
source_payload_size_decimal: provenance.source_payload_size_bytes().map(|value| return value.to_string()),
|
||||
};
|
||||
}
|
||||
|
||||
fn account_observation_row_dto(summary: &ksp_store_lib::RawAccountObservationSummary) -> crate::StoreAccountObservationRowDto {
|
||||
return crate::StoreAccountObservationRowDto {
|
||||
is_startup: summary.is_startup(),
|
||||
observation_key: encode_hex(summary.observation_key().as_bytes()),
|
||||
provenance: observation_provenance_dto(summary.provenance()),
|
||||
transaction_signature: summary.transaction_signature().map(|value| return encode_hex(value.as_bytes())),
|
||||
write_version_decimal: summary.write_version().map(|value| return value.to_string()),
|
||||
};
|
||||
}
|
||||
|
||||
fn transaction_observation_row_dto(summary: &ksp_store_lib::RawTransactionObservationSummary) -> crate::StoreTransactionObservationRowDto {
|
||||
return crate::StoreTransactionObservationRowDto {
|
||||
observation_key: encode_hex(summary.observation_key().as_bytes()),
|
||||
provenance: observation_provenance_dto(summary.provenance()),
|
||||
};
|
||||
}
|
||||
|
||||
fn acquisition_origin_code(origin: ksp_store_lib::RawAcquisitionOrigin) -> &'static str {
|
||||
return match origin {
|
||||
ksp_store_lib::RawAcquisitionOrigin::Backfill => "backfill",
|
||||
ksp_store_lib::RawAcquisitionOrigin::Import => "import",
|
||||
ksp_store_lib::RawAcquisitionOrigin::Live => "live",
|
||||
ksp_store_lib::RawAcquisitionOrigin::Repair => "repair",
|
||||
ksp_store_lib::RawAcquisitionOrigin::Replay => "replay",
|
||||
_ => "unknown",
|
||||
};
|
||||
}
|
||||
|
||||
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 {
|
||||
let page = inspection_page_from_request(request.offset.as_str(), request.limit);
|
||||
let page = match page {
|
||||
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()),
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let slot_min = parse_optional_decimal_u64(request.slot_min.as_deref());
|
||||
let slot_min = match slot_min {
|
||||
@@ -421,7 +608,6 @@ fn transaction_query_from_request(
|
||||
_ => 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));
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user