159 lines
7.8 KiB
Rust
159 lines
7.8 KiB
Rust
// file: crates/ksp-offchain-transport-lib/src/market_price_adapter.rs
|
|
// version: 2
|
|
|
|
//! Shared market-price adapter mechanics layered over crate-wide HTTP primitives.
|
|
|
|
/// Applies one non-blocking market-price admission decision and maps deferral to a stable KSP error.
|
|
pub(crate) fn admit_request(provider: &'static str, admission: &crate::HttpAdmissionController) -> ksp_core_lib::Result<()> {
|
|
return match admission.try_admit() {
|
|
crate::HttpAdmissionDecision::Ready => std::result::Result::Ok(()),
|
|
crate::HttpAdmissionDecision::Deferred(delay) => std::result::Result::Err(
|
|
ksp_core_lib::Error::new(crate::ERROR_CODE_HTTP_ADMISSION_DEFERRED, "Off-chain provider request is locally deferred")
|
|
.with_context("provider", provider)
|
|
.with_context("retry_after_millis", duration_millis_u64(delay).to_string()),
|
|
),
|
|
};
|
|
}
|
|
|
|
/// Captures the current UTC wall clock as a bounded market-price timestamp.
|
|
pub(crate) fn current_timestamp() -> ksp_core_lib::Result<crate::MarketPriceTimestamp> {
|
|
let duration = match std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => {
|
|
return std::result::Result::Err(
|
|
ksp_core_lib::Error::new(crate::ERROR_CODE_MARKET_PRICE_OBSERVATION_INVALID, "System clock cannot produce a market-price timestamp")
|
|
.with_source(error),
|
|
);
|
|
},
|
|
};
|
|
let millis = match u64::try_from(duration.as_millis()) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => {
|
|
return std::result::Result::Err(
|
|
ksp_core_lib::Error::new(crate::ERROR_CODE_MARKET_PRICE_OBSERVATION_INVALID, "System clock exceeds market-price timestamp bounds")
|
|
.with_source(error),
|
|
);
|
|
},
|
|
};
|
|
return std::result::Result::Ok(crate::MarketPriceTimestamp::from_unix_millis(millis));
|
|
}
|
|
|
|
/// Executes one provider GET and feeds any HTTP 429 cooldown back into the provider admission controller.
|
|
pub(crate) async fn get_json(
|
|
http: &crate::HttpRestClient,
|
|
admission: &crate::HttpAdmissionController,
|
|
provider: &'static str,
|
|
request: crate::HttpGetRequest,
|
|
) -> ksp_core_lib::Result<crate::HttpJsonDocument> {
|
|
let result = http.get_json(provider, "sol_usd", request).await;
|
|
if let std::result::Result::Err(error) = &result
|
|
&& error.code() == crate::ERROR_CODE_HTTP_RATE_LIMITED
|
|
{
|
|
let retry_after = retry_after_from_error(error);
|
|
admission.record_rate_limited(retry_after);
|
|
}
|
|
return result;
|
|
}
|
|
|
|
/// Builds one safe provider-response contract error without copying remote payload data.
|
|
pub(crate) fn invalid_provider_response(provider: &'static str, field: &'static str) -> ksp_core_lib::Error {
|
|
ksp_logging_lib::warn!(target: crate::TRACING_TARGET, provider = provider, field = field, "rejected invalid market-price provider response");
|
|
return ksp_core_lib::Error::new(crate::ERROR_CODE_MARKET_PRICE_PROVIDER_RESPONSE_INVALID, "Off-chain provider returned an invalid market-price response")
|
|
.with_context("provider", provider)
|
|
.with_context("field", field);
|
|
}
|
|
|
|
/// Builds one safe provider-response contract error and attaches a parser source that contains no remote payload copy.
|
|
pub(crate) fn invalid_provider_response_with_source<E>(provider: &'static str, field: &'static str, source: E) -> ksp_core_lib::Error
|
|
where
|
|
E: std::error::Error + std::marker::Send + std::marker::Sync + 'static,
|
|
{
|
|
return invalid_provider_response(provider, field).with_source(source);
|
|
}
|
|
|
|
/// Converts one provider RFC 3339 timestamp into the public millisecond timestamp contract.
|
|
pub(crate) fn market_price_timestamp_from_rfc3339(source: &str) -> std::option::Option<crate::MarketPriceTimestamp> {
|
|
let parsed = match chrono::DateTime::parse_from_rfc3339(source) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let millis = parsed.timestamp_millis();
|
|
if millis < 0 {
|
|
return std::option::Option::None;
|
|
}
|
|
return match u64::try_from(millis) {
|
|
std::result::Result::Ok(value) => std::option::Option::Some(crate::MarketPriceTimestamp::from_unix_millis(value)),
|
|
std::result::Result::Err(_) => std::option::Option::None,
|
|
};
|
|
}
|
|
|
|
/// Converts whole Unix seconds into the public millisecond timestamp contract with overflow checking.
|
|
pub(crate) fn market_price_timestamp_from_unix_seconds(seconds: u64) -> std::option::Option<crate::MarketPriceTimestamp> {
|
|
return seconds.checked_mul(1_000).map(crate::MarketPriceTimestamp::from_unix_millis);
|
|
}
|
|
|
|
/// Builds the stable error returned when a disabled provider is invoked directly.
|
|
pub(crate) fn provider_disabled_error(provider: &'static str) -> ksp_core_lib::Error {
|
|
return ksp_core_lib::Error::new(crate::ERROR_CODE_MARKET_PRICE_PROVIDER_DISABLED, "Off-chain market-price provider is disabled")
|
|
.with_context("provider", provider);
|
|
}
|
|
|
|
/// Builds hardened HTTP and admission runtime primitives from one provider-neutral rate-limit descriptor.
|
|
pub(crate) fn provider_http_runtime(
|
|
rate_limit: crate::MarketPriceProviderRateLimit,
|
|
) -> ksp_core_lib::Result<(crate::HttpRestClient, crate::HttpAdmissionController)> {
|
|
let policy = match rate_limit.kind() {
|
|
crate::MarketPriceProviderRateLimitKind::Dynamic => crate::HttpAdmissionPolicy::Dynamic,
|
|
crate::MarketPriceProviderRateLimitKind::Fixed => {
|
|
let requests = match rate_limit.requests() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(invalid_rate_limit_bridge()),
|
|
};
|
|
let window_seconds = match rate_limit.window_seconds() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return std::result::Result::Err(invalid_rate_limit_bridge()),
|
|
};
|
|
match crate::HttpAdmissionPolicy::fixed(requests, std::time::Duration::from_secs(u64::from(window_seconds)), rate_limit.burst()) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
}
|
|
},
|
|
};
|
|
let admission = match crate::HttpAdmissionController::new(policy, std::option::Option::None) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
let http = match crate::HttpRestClient::new(crate::HttpClientSettings::default()) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
};
|
|
return std::result::Result::Ok((http, admission));
|
|
}
|
|
|
|
fn duration_millis_u64(duration: std::time::Duration) -> u64 {
|
|
return match u64::try_from(duration.as_millis()) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => u64::MAX,
|
|
};
|
|
}
|
|
|
|
fn invalid_rate_limit_bridge() -> ksp_core_lib::Error {
|
|
return ksp_core_lib::Error::new(crate::ERROR_CODE_HTTP_RATE_LIMIT_INVALID, "Market-price rate-limit descriptor cannot map to HTTP admission policy");
|
|
}
|
|
|
|
fn retry_after_from_error(error: &ksp_core_lib::Error) -> std::option::Option<std::time::Duration> {
|
|
for context in error.context() {
|
|
if context.key() == "retry_after_seconds" {
|
|
let seconds = match context.value().parse::<u64>() {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
return std::option::Option::Some(std::time::Duration::from_secs(seconds));
|
|
}
|
|
}
|
|
return std::option::Option::None;
|
|
}
|
|
#[cfg(test)]
|
|
#[path = "../unit_tests/market_price_adapter.rs"]
|
|
mod tests;
|