v0.3.10-pre.003
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
# file: crates/ksp-job-backfill-lib/Cargo.toml
|
||||
# version: 4
|
||||
# version: 5
|
||||
|
||||
[package]
|
||||
name = "ksp-job-backfill-lib"
|
||||
@@ -13,6 +13,7 @@ ksp-core-lib = { path = "../ksp-core-lib" }
|
||||
ksp-job-api = { path = "../ksp-job-api" }
|
||||
ksp-logging-lib = { path = "../ksp-logging-lib" }
|
||||
ksp-onchain-transport-lib = { path = "../ksp-onchain-transport-lib" }
|
||||
ksp-raw-transaction-lib = { path = "../ksp-raw-transaction-lib" }
|
||||
ksp-store-lib = { path = "../ksp-store-lib", default-features = false }
|
||||
serde_json.workspace = true
|
||||
sha2.workspace = true
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
<!-- file: crates/ksp-job-backfill-lib/README.md -->
|
||||
<!-- version: 1 -->
|
||||
<!-- version: 2 -->
|
||||
|
||||
# ksp-job-backfill-lib
|
||||
|
||||
`ksp-job-backfill-lib` implémente le premier job historique concret de KSP : un backfill borné de transactions Solana vers la couche RAW durable.
|
||||
|
||||
La crate compose les contrats Job, Transport observé et Store backend-neutral sans posséder leurs politiques internes. Elle couvre l'admission, la découverte historique, l'hydratation `getTransaction`, la conversion RAW v1, la persistance atomique, la concurrence bornée, la frontier contiguë, le checkpoint caller-owned, l'annulation coopérative et les snapshots latest-value.
|
||||
La crate compose les contrats Job, Transport observé, common RAW et Store backend-neutral sans posséder leurs politiques internes. Elle couvre l'admission, la découverte historique, l'hydratation `getTransaction`, l'adaptation vers `ksp-raw-transaction-lib`, la persistance atomique, la concurrence bornée, la frontier contiguë, le checkpoint caller-owned, l'annulation coopérative et les snapshots latest-value.
|
||||
|
||||
## Identité
|
||||
|
||||
@@ -94,6 +94,7 @@ ksp-core-lib
|
||||
ksp-job-api
|
||||
ksp-logging-lib
|
||||
ksp-onchain-transport-lib
|
||||
ksp-raw-transaction-lib
|
||||
ksp-store-lib (default-features = false)
|
||||
futures-util
|
||||
serde_json
|
||||
|
||||
@@ -1,18 +1,13 @@
|
||||
// file: crates/ksp-job-backfill-lib/src/conversion.rs
|
||||
// version: 4
|
||||
// version: 5
|
||||
|
||||
use sha2::Digest; // rust-rules: trait-import
|
||||
|
||||
/// KSP-owned source-independent RAW transaction format identifier produced by this Backfill vertical.
|
||||
pub const RAW_TRANSACTION_FORMAT_ID: &str = "ksp.solana.raw_transaction";
|
||||
/// Initial KSP-owned RAW transaction format version produced by this Backfill vertical.
|
||||
pub const RAW_TRANSACTION_FORMAT_VERSION: u32 = 1;
|
||||
|
||||
const RAW_TRANSACTION_METHOD_CODE: &str = "getTransaction";
|
||||
const RAW_TRANSACTION_OBSERVATION_CONTRACT_VERSION: u32 = 1;
|
||||
const RAW_TRANSACTION_PROTOCOL_CODE: &str = "solana.http.json_rpc";
|
||||
|
||||
/// Complete in-memory RAW transaction acquisition ready for the later Store persistence tranche.
|
||||
/// Complete in-memory RAW transaction acquisition ready for Store persistence.
|
||||
#[derive(Debug)]
|
||||
pub struct BackfillRawAcquisition {
|
||||
inner: Box<BackfillRawAcquisitionInner>,
|
||||
@@ -20,35 +15,34 @@ pub struct BackfillRawAcquisition {
|
||||
|
||||
#[derive(Debug)]
|
||||
struct BackfillRawAcquisitionInner {
|
||||
transaction: ksp_store_lib::RawTransaction,
|
||||
observation: ksp_store_lib::RawTransactionObservation,
|
||||
acquisition: ksp_raw_transaction_lib::RawTransactionAcquisition,
|
||||
}
|
||||
|
||||
impl crate::BackfillRawAcquisition {
|
||||
/// Returns the canonical RAW transaction produced from the typed Transport response.
|
||||
/// Returns the canonical RAW transaction produced through the source-neutral common layer.
|
||||
#[must_use]
|
||||
pub const fn transaction(&self) -> &ksp_store_lib::RawTransaction {
|
||||
return &self.inner.transaction;
|
||||
pub fn transaction(&self) -> &ksp_store_lib::RawTransaction {
|
||||
return self.inner.acquisition.transaction();
|
||||
}
|
||||
|
||||
/// Returns the acquisition observation whose provenance records the actual successful endpoint.
|
||||
/// Returns the Backfill-owned acquisition observation linked by the common assembly layer.
|
||||
#[must_use]
|
||||
pub const fn observation(&self) -> &ksp_store_lib::RawTransactionObservation {
|
||||
return &self.inner.observation;
|
||||
pub fn observation(&self) -> &ksp_store_lib::RawTransactionObservation {
|
||||
return self.inner.acquisition.observation();
|
||||
}
|
||||
|
||||
/// Consumes the in-memory acquisition into the canonical transaction and its observation.
|
||||
/// Consumes the in-memory acquisition into the canonical transaction and its Backfill-owned observation.
|
||||
#[must_use]
|
||||
pub fn into_parts(self) -> (ksp_store_lib::RawTransaction, ksp_store_lib::RawTransactionObservation) {
|
||||
let inner = *self.inner;
|
||||
return (inner.transaction, inner.observation);
|
||||
return inner.acquisition.into_parts();
|
||||
}
|
||||
}
|
||||
|
||||
/// Result of hydrating one deterministic Backfill candidate through observed `getTransaction`.
|
||||
#[derive(Debug)]
|
||||
pub enum BackfillHydrationOutcome {
|
||||
/// The RPC returned one complete transaction and conversion produced canonical RAW plus provenance.
|
||||
/// The RPC returned one complete transaction and common conversion produced canonical RAW plus Backfill provenance.
|
||||
Available(crate::BackfillRawAcquisition),
|
||||
/// The RPC returned JSON `null`; only the canonical transaction identity exists and no provenance is fabricated.
|
||||
Missing(ksp_store_lib::RawTransactionReference),
|
||||
@@ -71,11 +65,12 @@ impl crate::BackfillHydrationOutcome {
|
||||
}
|
||||
}
|
||||
|
||||
/// Hydrates one candidate with the typed observed Transport path and converts it to canonical RAW v1.
|
||||
/// Hydrates one candidate with the typed observed Transport path and converts it through the common RAW v1 layer.
|
||||
///
|
||||
/// The caller supplies the local receipt timestamp because wall-clock ownership remains outside this
|
||||
/// pure conversion tranche. Transport retains endpoint selection and retry. This function never
|
||||
/// persists to Store; persistence begins in `pre.007`.
|
||||
/// pure conversion path. Transport retains endpoint selection and retry. Backfill retains campaign
|
||||
/// provenance and observation-key ownership; the common crate owns signature parsing, RAW v1
|
||||
/// canonicalization, content hashing and transaction/observation assembly.
|
||||
pub async fn hydrate_backfill_candidate(
|
||||
transport: &ksp_onchain_transport_lib::HttpTransportPool,
|
||||
request: &crate::BackfillRequest,
|
||||
@@ -105,12 +100,12 @@ pub async fn hydrate_backfill_candidate(
|
||||
std::option::Option::None => return std::result::Result::Ok(crate::BackfillHydrationOutcome::Missing(reference)),
|
||||
};
|
||||
let fields = CanonicalTransactionFields {
|
||||
slot: transaction.slot(),
|
||||
block_time: transaction.block_time(),
|
||||
transaction: transaction.transaction(),
|
||||
meta: transaction.meta(),
|
||||
version: transaction.version(),
|
||||
slot: transaction.slot(),
|
||||
transaction: transaction.transaction(),
|
||||
transaction_index: transaction.transaction_index(),
|
||||
version: transaction.version(),
|
||||
};
|
||||
let acquisition = convert_available_fields(request, reference, fields, provider.as_str(), endpoint.as_str(), received_at);
|
||||
return match acquisition {
|
||||
@@ -119,49 +114,22 @@ pub async fn hydrate_backfill_candidate(
|
||||
};
|
||||
}
|
||||
|
||||
/// Decodes one validated Base58 signature to exactly 64 canonical bytes without a Solana SDK dependency.
|
||||
/// Decodes one validated Base58 signature through the source-neutral common parser while preserving Backfill error semantics.
|
||||
pub(crate) fn decode_backfill_signature(signature: &crate::BackfillSignature) -> ksp_core_lib::Result<ksp_store_lib::RawTransactionSignature> {
|
||||
let text = signature.as_str().as_bytes();
|
||||
let mut decoded = [0_u8; 64];
|
||||
let mut leading_zeroes = 0_usize;
|
||||
for byte in text {
|
||||
if *byte != b'1' {
|
||||
break;
|
||||
}
|
||||
leading_zeroes += 1;
|
||||
}
|
||||
for byte in text {
|
||||
let digit = match base58_digit(*byte) {
|
||||
std::option::Option::Some(digit) => digit,
|
||||
std::option::Option::None => return std::result::Result::Err(conversion_error("signature")),
|
||||
};
|
||||
let mut carry = u32::from(digit);
|
||||
for output in decoded.iter_mut().rev() {
|
||||
let value = (u32::from(*output) * 58) + carry;
|
||||
*output = (value & 0xff) as u8;
|
||||
carry = value >> 8;
|
||||
}
|
||||
if carry != 0 {
|
||||
return std::result::Result::Err(conversion_error("signature"));
|
||||
}
|
||||
}
|
||||
let significant_len = match decoded.iter().position(|byte| return *byte != 0) {
|
||||
std::option::Option::Some(index) => decoded.len() - index,
|
||||
std::option::Option::None => 0,
|
||||
let parsed = ksp_raw_transaction_lib::parse_raw_transaction_signature(signature.as_str());
|
||||
return match parsed {
|
||||
std::result::Result::Ok(signature) => std::result::Result::Ok(signature),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(conversion_error("signature")),
|
||||
};
|
||||
if leading_zeroes + significant_len != decoded.len() {
|
||||
return std::result::Result::Err(conversion_error("signature"));
|
||||
}
|
||||
return std::result::Result::Ok(ksp_store_lib::RawTransactionSignature::new(decoded));
|
||||
}
|
||||
|
||||
struct CanonicalTransactionFields<'a> {
|
||||
slot: u64,
|
||||
block_time: std::option::Option<i64>,
|
||||
transaction: &'a ksp_onchain_transport_lib::SolanaEncodedTransaction,
|
||||
meta: &'a ksp_onchain_transport_lib::SolanaWireField<serde_json::Value>,
|
||||
version: &'a ksp_onchain_transport_lib::SolanaWireField<ksp_onchain_transport_lib::SolanaTransactionVersion>,
|
||||
slot: u64,
|
||||
transaction: &'a ksp_onchain_transport_lib::SolanaEncodedTransaction,
|
||||
transaction_index: &'a ksp_onchain_transport_lib::SolanaWireField<u32>,
|
||||
version: &'a ksp_onchain_transport_lib::SolanaWireField<ksp_onchain_transport_lib::SolanaTransactionVersion>,
|
||||
}
|
||||
|
||||
fn canonical_reference(request: &crate::BackfillRequest, candidate: &crate::BackfillCandidate) -> ksp_core_lib::Result<ksp_store_lib::RawTransactionReference> {
|
||||
@@ -184,30 +152,9 @@ fn convert_available_fields(
|
||||
endpoint: &str,
|
||||
received_at: ksp_store_lib::RawTimestamp,
|
||||
) -> ksp_core_lib::Result<crate::BackfillRawAcquisition> {
|
||||
let block_time = convert_block_time(fields.block_time);
|
||||
let block_time = match block_time {
|
||||
std::result::Result::Ok(block_time) => block_time,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let bytes = canonical_payload_bytes(&fields);
|
||||
let bytes = match bytes {
|
||||
std::result::Result::Ok(bytes) => bytes,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let hash: [u8; 32] = sha2::Sha256::digest(bytes.as_slice()).into();
|
||||
let format_id = ksp_store_lib::RawFormatId::new(crate::RAW_TRANSACTION_FORMAT_ID);
|
||||
let format_id = match format_id {
|
||||
std::result::Result::Ok(format_id) => format_id,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(conversion_error("payload.format_id")),
|
||||
};
|
||||
let payload = ksp_store_lib::RawPayload::try_new(
|
||||
format_id,
|
||||
crate::RAW_TRANSACTION_FORMAT_VERSION,
|
||||
bytes.into_boxed_slice(),
|
||||
ksp_store_lib::RawContentHash::new(hash),
|
||||
);
|
||||
let payload = match payload {
|
||||
std::result::Result::Ok(payload) => payload,
|
||||
let transaction = canonical_transaction(&reference, fields);
|
||||
let transaction = match transaction {
|
||||
std::result::Result::Ok(transaction) => transaction,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let provenance = build_provenance(request, provider, endpoint, received_at);
|
||||
@@ -215,182 +162,60 @@ fn convert_available_fields(
|
||||
std::result::Result::Ok(provenance) => provenance,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let observation_key = observation_key(request, &reference, provider, endpoint);
|
||||
let transaction = ksp_store_lib::RawTransaction::new(reference.clone(), fields.slot, block_time, payload);
|
||||
let observation = ksp_store_lib::RawTransactionObservation::new(observation_key, reference, provenance);
|
||||
return std::result::Result::Ok(crate::BackfillRawAcquisition { inner: Box::new(BackfillRawAcquisitionInner { transaction, observation }) });
|
||||
let observation_key = observation_key(request, transaction.reference(), provider, endpoint);
|
||||
let acquisition = ksp_raw_transaction_lib::assemble_raw_transaction_acquisition(transaction, observation_key, provenance);
|
||||
return std::result::Result::Ok(crate::BackfillRawAcquisition { inner: Box::new(BackfillRawAcquisitionInner { acquisition }) });
|
||||
}
|
||||
|
||||
fn convert_block_time(value: std::option::Option<i64>) -> ksp_core_lib::Result<std::option::Option<ksp_store_lib::RawTimestamp>> {
|
||||
let seconds = match value {
|
||||
std::option::Option::Some(seconds) => seconds,
|
||||
std::option::Option::None => return std::result::Result::Ok(std::option::Option::None),
|
||||
};
|
||||
let seconds = match u64::try_from(seconds) {
|
||||
std::result::Result::Ok(seconds) => seconds,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(conversion_error("block_time")),
|
||||
};
|
||||
let millis = match seconds.checked_mul(1_000) {
|
||||
std::option::Option::Some(millis) => millis,
|
||||
std::option::Option::None => return std::result::Result::Err(conversion_error("block_time")),
|
||||
};
|
||||
let timestamp = ksp_store_lib::RawTimestamp::from_unix_millis(millis);
|
||||
return match timestamp {
|
||||
std::result::Result::Ok(timestamp) => std::result::Result::Ok(std::option::Option::Some(timestamp)),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(conversion_error("block_time")),
|
||||
};
|
||||
}
|
||||
|
||||
fn canonical_payload_bytes(fields: &CanonicalTransactionFields<'_>) -> ksp_core_lib::Result<std::vec::Vec<u8>> {
|
||||
let (transaction_data, transaction_encoding) = match fields.transaction {
|
||||
ksp_onchain_transport_lib::SolanaEncodedTransaction::Binary { data, encoding } => {
|
||||
if *encoding != ksp_onchain_transport_lib::SolanaTransactionBinaryEncoding::Base64 {
|
||||
return std::result::Result::Err(conversion_error("transaction.encoding"));
|
||||
}
|
||||
(data.as_str(), "base64")
|
||||
fn canonical_transaction(
|
||||
reference: &ksp_store_lib::RawTransactionReference,
|
||||
fields: CanonicalTransactionFields<'_>,
|
||||
) -> ksp_core_lib::Result<ksp_store_lib::RawTransaction> {
|
||||
let transaction_data = match fields.transaction {
|
||||
ksp_onchain_transport_lib::SolanaEncodedTransaction::Binary { data, encoding }
|
||||
if *encoding == ksp_onchain_transport_lib::SolanaTransactionBinaryEncoding::Base64 =>
|
||||
{
|
||||
data.clone()
|
||||
},
|
||||
ksp_onchain_transport_lib::SolanaEncodedTransaction::LegacyBinary(_) | ksp_onchain_transport_lib::SolanaEncodedTransaction::Json(_) => {
|
||||
ksp_onchain_transport_lib::SolanaEncodedTransaction::Binary { .. }
|
||||
| ksp_onchain_transport_lib::SolanaEncodedTransaction::LegacyBinary(_)
|
||||
| ksp_onchain_transport_lib::SolanaEncodedTransaction::Json(_) => {
|
||||
return std::result::Result::Err(conversion_error("transaction.encoding"));
|
||||
},
|
||||
};
|
||||
let mut bytes = std::vec::Vec::new();
|
||||
bytes.extend_from_slice(b"{\"transaction\":[");
|
||||
if let std::result::Result::Err(error) = append_json_string(&mut bytes, transaction_data) {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
bytes.push(b',');
|
||||
if let std::result::Result::Err(error) = append_json_string(&mut bytes, transaction_encoding) {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
bytes.push(b']');
|
||||
if let std::result::Result::Err(error) = append_wire_value(&mut bytes, "meta", fields.meta, |output, value| return append_canonical_json(output, value)) {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
let version_result = append_wire_value(&mut bytes, "version", fields.version, |output, value| {
|
||||
let meta = map_wire_field(fields.meta, |value| return value.clone());
|
||||
let version = map_wire_field(fields.version, |value| {
|
||||
return match value {
|
||||
ksp_onchain_transport_lib::SolanaTransactionVersion::Legacy => append_json_string(output, "legacy"),
|
||||
ksp_onchain_transport_lib::SolanaTransactionVersion::Number(number) => {
|
||||
output.extend_from_slice(number.to_string().as_bytes());
|
||||
std::result::Result::Ok(())
|
||||
},
|
||||
ksp_onchain_transport_lib::SolanaTransactionVersion::Legacy => ksp_raw_transaction_lib::RawTransactionVersion::Legacy,
|
||||
ksp_onchain_transport_lib::SolanaTransactionVersion::Number(number) => ksp_raw_transaction_lib::RawTransactionVersion::Number(*number),
|
||||
};
|
||||
});
|
||||
if let std::result::Result::Err(error) = version_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
let transaction_index_result = append_wire_value(&mut bytes, "transactionIndex", fields.transaction_index, |output, value| {
|
||||
output.extend_from_slice(value.to_string().as_bytes());
|
||||
return std::result::Result::Ok(());
|
||||
});
|
||||
if let std::result::Result::Err(error) = transaction_index_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
bytes.push(b'}');
|
||||
return std::result::Result::Ok(bytes);
|
||||
let transaction_index = map_wire_field(fields.transaction_index, |value| return *value);
|
||||
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64(
|
||||
reference.network().clone(),
|
||||
reference.signature(),
|
||||
fields.slot,
|
||||
fields.block_time,
|
||||
transaction_data,
|
||||
meta,
|
||||
version,
|
||||
transaction_index,
|
||||
);
|
||||
let canonical = ksp_raw_transaction_lib::canonicalize_raw_transaction(material);
|
||||
return match canonical {
|
||||
std::result::Result::Ok(transaction) => std::result::Result::Ok(transaction),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(conversion_error("common.canonicalization")),
|
||||
};
|
||||
}
|
||||
|
||||
fn append_wire_value<T, F>(
|
||||
output: &mut std::vec::Vec<u8>,
|
||||
key: &str,
|
||||
field: &ksp_onchain_transport_lib::SolanaWireField<T>,
|
||||
mut append_value: F,
|
||||
) -> ksp_core_lib::Result<()>
|
||||
fn map_wire_field<T, U, F>(field: &ksp_onchain_transport_lib::SolanaWireField<T>, mut map_value: F) -> ksp_raw_transaction_lib::RawTransactionWireField<U>
|
||||
where
|
||||
F: FnMut(&mut std::vec::Vec<u8>, &T) -> ksp_core_lib::Result<()>,
|
||||
F: FnMut(&T) -> U,
|
||||
{
|
||||
return match field {
|
||||
ksp_onchain_transport_lib::SolanaWireField::Omitted => std::result::Result::Ok(()),
|
||||
ksp_onchain_transport_lib::SolanaWireField::Null => {
|
||||
output.push(b',');
|
||||
let key_result = append_json_string(output, key);
|
||||
if let std::result::Result::Err(error) = key_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
output.extend_from_slice(b":null");
|
||||
std::result::Result::Ok(())
|
||||
},
|
||||
ksp_onchain_transport_lib::SolanaWireField::Value(value) => {
|
||||
output.push(b',');
|
||||
let key_result = append_json_string(output, key);
|
||||
if let std::result::Result::Err(error) = key_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
output.push(b':');
|
||||
append_value(output, value)
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
fn append_canonical_json(output: &mut std::vec::Vec<u8>, value: &serde_json::Value) -> ksp_core_lib::Result<()> {
|
||||
return match value {
|
||||
serde_json::Value::Null => {
|
||||
output.extend_from_slice(b"null");
|
||||
std::result::Result::Ok(())
|
||||
},
|
||||
serde_json::Value::Bool(value) => {
|
||||
if *value {
|
||||
output.extend_from_slice(b"true");
|
||||
} else {
|
||||
output.extend_from_slice(b"false");
|
||||
}
|
||||
std::result::Result::Ok(())
|
||||
},
|
||||
serde_json::Value::Number(value) => {
|
||||
output.extend_from_slice(value.to_string().as_bytes());
|
||||
std::result::Result::Ok(())
|
||||
},
|
||||
serde_json::Value::String(value) => append_json_string(output, value.as_str()),
|
||||
serde_json::Value::Array(values) => {
|
||||
output.push(b'[');
|
||||
for (index, item) in values.iter().enumerate() {
|
||||
if index != 0 {
|
||||
output.push(b',');
|
||||
}
|
||||
let item_result = append_canonical_json(output, item);
|
||||
if let std::result::Result::Err(error) = item_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
}
|
||||
output.push(b']');
|
||||
std::result::Result::Ok(())
|
||||
},
|
||||
serde_json::Value::Object(values) => {
|
||||
output.push(b'{');
|
||||
let mut keys = values.keys().collect::<std::vec::Vec<_>>();
|
||||
keys.sort_unstable();
|
||||
for (index, key) in keys.iter().enumerate() {
|
||||
if index != 0 {
|
||||
output.push(b',');
|
||||
}
|
||||
let key_result = append_json_string(output, key.as_str());
|
||||
if let std::result::Result::Err(error) = key_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
output.push(b':');
|
||||
let item = values.get(key.as_str());
|
||||
let item = match item {
|
||||
std::option::Option::Some(item) => item,
|
||||
std::option::Option::None => return std::result::Result::Err(conversion_error("payload.meta")),
|
||||
};
|
||||
let item_result = append_canonical_json(output, item);
|
||||
if let std::result::Result::Err(error) = item_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
}
|
||||
output.push(b'}');
|
||||
std::result::Result::Ok(())
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
fn append_json_string(output: &mut std::vec::Vec<u8>, value: &str) -> ksp_core_lib::Result<()> {
|
||||
let encoded = serde_json::to_vec(value);
|
||||
return match encoded {
|
||||
std::result::Result::Ok(encoded) => {
|
||||
output.extend_from_slice(encoded.as_slice());
|
||||
std::result::Result::Ok(())
|
||||
},
|
||||
std::result::Result::Err(_) => std::result::Result::Err(conversion_error("payload.json")),
|
||||
ksp_onchain_transport_lib::SolanaWireField::Omitted => ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
||||
ksp_onchain_transport_lib::SolanaWireField::Null => ksp_raw_transaction_lib::RawTransactionWireField::Null,
|
||||
ksp_onchain_transport_lib::SolanaWireField::Value(value) => ksp_raw_transaction_lib::RawTransactionWireField::Value(map_value(value)),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -463,18 +288,6 @@ fn hash_bytes(hasher: &mut sha2::Sha256, value: &[u8]) {
|
||||
hasher.update(value);
|
||||
}
|
||||
|
||||
fn base58_digit(byte: u8) -> std::option::Option<u8> {
|
||||
return match byte {
|
||||
b'1'..=b'9' => std::option::Option::Some(byte - b'1'),
|
||||
b'A'..=b'H' => std::option::Option::Some((byte - b'A') + 9),
|
||||
b'J'..=b'N' => std::option::Option::Some((byte - b'J') + 17),
|
||||
b'P'..=b'Z' => std::option::Option::Some((byte - b'P') + 22),
|
||||
b'a'..=b'k' => std::option::Option::Some((byte - b'a') + 33),
|
||||
b'm'..=b'z' => std::option::Option::Some((byte - b'm') + 44),
|
||||
_ => std::option::Option::None,
|
||||
};
|
||||
}
|
||||
|
||||
fn conversion_error(field: &'static str) -> ksp_core_lib::Error {
|
||||
return ksp_core_lib::Error::new(crate::ERROR_CODE_BACKFILL_RAW_CONVERSION_INVALID, "invalid deterministic Backfill RAW conversion")
|
||||
.with_context("field", field);
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-job-backfill-lib/src/lib.rs
|
||||
// version: 6
|
||||
// version: 7
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -30,10 +30,6 @@ pub use self::checkpoint::BackfillCheckpoint;
|
||||
pub use self::conversion::BackfillHydrationOutcome;
|
||||
/// Complete in-memory RAW transaction acquisition ready for later Store persistence.
|
||||
pub use self::conversion::BackfillRawAcquisition;
|
||||
/// KSP-owned source-independent RAW transaction format identifier produced by this Backfill vertical.
|
||||
pub use self::conversion::RAW_TRANSACTION_FORMAT_ID;
|
||||
/// Initial KSP-owned RAW transaction format version produced by this Backfill vertical.
|
||||
pub use self::conversion::RAW_TRANSACTION_FORMAT_VERSION;
|
||||
/// Hydrates one candidate through observed Transport and converts a non-null response to canonical RAW v1.
|
||||
pub use self::conversion::hydrate_backfill_candidate;
|
||||
/// One deterministic transaction candidate produced by bounded discovery.
|
||||
@@ -112,6 +108,10 @@ pub use self::runtime::BackfillJobRuntime;
|
||||
pub use self::runtime::BackfillJobSnapshot;
|
||||
/// Cloneable runtime-neutral-facing latest-value source for concrete Backfill snapshots.
|
||||
pub use self::runtime::BackfillSnapshotSource;
|
||||
/// KSP-owned source-independent RAW transaction format identifier re-exported from the common RAW layer.
|
||||
pub use ksp_raw_transaction_lib::RAW_TRANSACTION_FORMAT_ID;
|
||||
/// Frozen RAW transaction format version re-exported from the common RAW layer.
|
||||
pub use ksp_raw_transaction_lib::RAW_TRANSACTION_FORMAT_VERSION;
|
||||
|
||||
/// Internal contiguous completion frontier used by bounded execution.
|
||||
pub(crate) use self::checkpoint::CompletionFrontier;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-job-backfill-lib/src/request.rs
|
||||
// version: 6
|
||||
// version: 7
|
||||
|
||||
use sha2::Digest; // rust-rules: trait-import
|
||||
|
||||
@@ -12,9 +12,9 @@ pub const MAX_BACKFILL_PAGES: usize = 10_000;
|
||||
/// Maximum page size admitted for one `getSignaturesForAddress` request.
|
||||
pub const MAX_BACKFILL_PAGE_SIZE: usize = 1_000;
|
||||
/// Maximum Base58 text length possible for one canonical 64-byte Solana signature.
|
||||
pub const MAX_BACKFILL_SIGNATURE_TEXT_BYTES: usize = 88;
|
||||
pub const MAX_BACKFILL_SIGNATURE_TEXT_BYTES: usize = ksp_raw_transaction_lib::MAX_RAW_TRANSACTION_SIGNATURE_TEXT_BYTES;
|
||||
/// Minimum Base58 text length possible for one canonical 64-byte Solana signature.
|
||||
pub const MIN_BACKFILL_SIGNATURE_TEXT_BYTES: usize = 64;
|
||||
pub const MIN_BACKFILL_SIGNATURE_TEXT_BYTES: usize = ksp_raw_transaction_lib::MIN_RAW_TRANSACTION_SIGNATURE_TEXT_BYTES;
|
||||
|
||||
/// Commitment levels intentionally admitted by the historical Backfill vertical.
|
||||
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
||||
|
||||
@@ -1,10 +1,10 @@
|
||||
// file: crates/ksp-job-backfill-lib/tests/dependency_boundary.rs
|
||||
// version: 6
|
||||
// version: 7
|
||||
|
||||
//! Dependency firewall canaries through the concrete cancellation and latest-value runtime tranche.
|
||||
|
||||
#[test]
|
||||
fn pre_009_manifest_uses_only_planned_ksp_edges_and_private_tokio_runtime() {
|
||||
fn pre_003_manifest_adds_common_raw_edge_without_crossing_store_or_runtime_boundaries() {
|
||||
let manifest = include_str!("../Cargo.toml");
|
||||
for required in [
|
||||
"futures-util = { workspace = true, features = [\"std\"] }",
|
||||
@@ -12,6 +12,7 @@ fn pre_009_manifest_uses_only_planned_ksp_edges_and_private_tokio_runtime() {
|
||||
"ksp-job-api = { path = \"../ksp-job-api\" }",
|
||||
"ksp-logging-lib = { path = \"../ksp-logging-lib\" }",
|
||||
"ksp-onchain-transport-lib = { path = \"../ksp-onchain-transport-lib\" }",
|
||||
"ksp-raw-transaction-lib = { path = \"../ksp-raw-transaction-lib\" }",
|
||||
"ksp-store-lib = { path = \"../ksp-store-lib\", default-features = false }",
|
||||
"serde_json.workspace = true",
|
||||
"sha2.workspace = true",
|
||||
@@ -56,7 +57,12 @@ fn pre_009_production_sources_keep_transport_store_and_scheduler_ownership_separ
|
||||
}
|
||||
}
|
||||
let conversion = include_str!("../src/conversion.rs");
|
||||
assert!(conversion.contains("serde_json::"));
|
||||
assert!(conversion.contains("ksp_raw_transaction_lib::canonicalize_raw_transaction"));
|
||||
assert!(conversion.contains("ksp_raw_transaction_lib::assemble_raw_transaction_acquisition"));
|
||||
assert!(conversion.contains("ksp_raw_transaction_lib::parse_raw_transaction_signature"));
|
||||
assert!(!conversion.contains("fn append_canonical_json"));
|
||||
assert!(!conversion.contains("fn base58_digit"));
|
||||
assert!(!conversion.contains("RawPayload::try_new"));
|
||||
assert!(conversion.contains("get_transaction_observed"));
|
||||
assert!(conversion.contains("SolanaTransactionEncoding::Base64"));
|
||||
assert!(conversion.contains("std::option::Option::Some(0)"));
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-job-backfill-lib/tests/hardening.rs
|
||||
// version: 4
|
||||
// version: 5
|
||||
|
||||
//! Adversarial, security, visibility and external-boundary hardening canaries for `pre.010`.
|
||||
|
||||
@@ -268,6 +268,7 @@ fn pre_010_manifest_dependency_surface_is_exact_and_backend_neutral() {
|
||||
"ksp-job-api",
|
||||
"ksp-logging-lib",
|
||||
"ksp-onchain-transport-lib",
|
||||
"ksp-raw-transaction-lib",
|
||||
"ksp-store-lib",
|
||||
"serde_json.workspace",
|
||||
"sha2.workspace",
|
||||
@@ -328,6 +329,7 @@ fn pre_010_lower_layers_have_no_dependency_return_to_job() {
|
||||
include_str!("../../ksp-core-lib/Cargo.toml"),
|
||||
include_str!("../../ksp-logging-lib/Cargo.toml"),
|
||||
include_str!("../../ksp-onchain-transport-lib/Cargo.toml"),
|
||||
include_str!("../../ksp-raw-transaction-lib/Cargo.toml"),
|
||||
include_str!("../../ksp-store-api/Cargo.toml"),
|
||||
include_str!("../../ksp-store-lib/Cargo.toml"),
|
||||
include_str!("../../ksp-store-postgres-lib/Cargo.toml"),
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-job-backfill-lib/unit_tests/conversion.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
fn signature_text() -> std::option::Option<crate::BackfillSignature> {
|
||||
return match crate::BackfillSignature::new("1".repeat(64)) {
|
||||
@@ -46,6 +46,14 @@ fn candidate(network: &str, signature: crate::BackfillSignature) -> std::option:
|
||||
return std::option::Option::Some(crate::BackfillCandidate::new(crate::BackfillCandidateIdentity::new(network, signature), std::option::Option::Some(42)));
|
||||
}
|
||||
|
||||
fn raw_reference(network: &str) -> std::option::Option<ksp_store_lib::RawTransactionReference> {
|
||||
let network = match ksp_store_lib::RawNetworkId::new(network) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::option::Option::None,
|
||||
};
|
||||
return std::option::Option::Some(ksp_store_lib::RawTransactionReference::new(network, ksp_store_lib::RawTransactionSignature::new([0_u8; 64])));
|
||||
}
|
||||
|
||||
fn received_at() -> std::option::Option<ksp_store_lib::RawTimestamp> {
|
||||
return match ksp_store_lib::RawTimestamp::from_unix_millis(1_700_000_001_000) {
|
||||
std::result::Result::Ok(value) => std::option::Option::Some(value),
|
||||
@@ -169,21 +177,26 @@ fn pre_006_wire_omission_and_null_produce_distinct_canonical_bytes() {
|
||||
let version_null = ksp_onchain_transport_lib::SolanaWireField::Null;
|
||||
let transaction_index_omitted = ksp_onchain_transport_lib::SolanaWireField::Omitted;
|
||||
let transaction_index_null = ksp_onchain_transport_lib::SolanaWireField::Null;
|
||||
let omitted = super::canonical_payload_bytes(&fields(&transaction, &meta_omitted, &version_omitted, &transaction_index_omitted, std::option::Option::None));
|
||||
let reference = match raw_reference("devnet") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let omitted =
|
||||
super::canonical_transaction(&reference, fields(&transaction, &meta_omitted, &version_omitted, &transaction_index_omitted, std::option::Option::None));
|
||||
assert!(omitted.is_ok());
|
||||
let omitted = match omitted {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return,
|
||||
};
|
||||
assert_eq!(omitted, b"{\"transaction\":[\"AQID\",\"base64\"]}");
|
||||
let nulls = super::canonical_payload_bytes(&fields(&transaction, &meta_null, &version_null, &transaction_index_null, std::option::Option::None));
|
||||
assert_eq!(omitted.payload().bytes(), b"{\"transaction\":[\"AQID\",\"base64\"]}");
|
||||
let nulls = super::canonical_transaction(&reference, fields(&transaction, &meta_null, &version_null, &transaction_index_null, std::option::Option::None));
|
||||
assert!(nulls.is_ok());
|
||||
let nulls = match nulls {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return,
|
||||
};
|
||||
assert_eq!(nulls, b"{\"transaction\":[\"AQID\",\"base64\"],\"meta\":null,\"version\":null,\"transactionIndex\":null}");
|
||||
assert_ne!(omitted, nulls);
|
||||
assert_eq!(nulls.payload().bytes(), b"{\"transaction\":[\"AQID\",\"base64\"],\"meta\":null,\"version\":null,\"transactionIndex\":null}");
|
||||
assert_ne!(omitted.payload().bytes(), nulls.payload().bytes());
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -198,8 +211,12 @@ fn pre_006_non_base64_transaction_shapes_are_rejected() {
|
||||
let meta = ksp_onchain_transport_lib::SolanaWireField::Omitted;
|
||||
let version = ksp_onchain_transport_lib::SolanaWireField::Omitted;
|
||||
let transaction_index = ksp_onchain_transport_lib::SolanaWireField::Omitted;
|
||||
let reference = match raw_reference("devnet") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
for transaction in [&base58, &legacy, &json] {
|
||||
let result = super::canonical_payload_bytes(&fields(transaction, &meta, &version, &transaction_index, std::option::Option::None));
|
||||
let result = super::canonical_transaction(&reference, fields(transaction, &meta, &version, &transaction_index, std::option::Option::None));
|
||||
assert!(result.is_err());
|
||||
if let std::result::Result::Err(error) = result {
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_BACKFILL_RAW_CONVERSION_INVALID);
|
||||
@@ -210,12 +227,26 @@ fn pre_006_non_base64_transaction_shapes_are_rejected() {
|
||||
|
||||
#[test]
|
||||
fn pre_006_negative_and_unrepresentable_block_times_are_terminal_conversion_errors() {
|
||||
let negative = super::convert_block_time(std::option::Option::Some(-1));
|
||||
let reference = match raw_reference("devnet") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return,
|
||||
};
|
||||
let transaction = ksp_onchain_transport_lib::SolanaEncodedTransaction::Binary {
|
||||
data: "AQID".to_owned(),
|
||||
encoding: ksp_onchain_transport_lib::SolanaTransactionBinaryEncoding::Base64,
|
||||
};
|
||||
let meta = ksp_onchain_transport_lib::SolanaWireField::Omitted;
|
||||
let version = ksp_onchain_transport_lib::SolanaWireField::Omitted;
|
||||
let transaction_index = ksp_onchain_transport_lib::SolanaWireField::Omitted;
|
||||
let negative = super::canonical_transaction(&reference, fields(&transaction, &meta, &version, &transaction_index, std::option::Option::Some(-1)));
|
||||
assert!(negative.is_err());
|
||||
let oversized = super::convert_block_time(std::option::Option::Some(i64::MAX));
|
||||
let oversized = super::canonical_transaction(&reference, fields(&transaction, &meta, &version, &transaction_index, std::option::Option::Some(i64::MAX)));
|
||||
assert!(oversized.is_err());
|
||||
let absent = super::convert_block_time(std::option::Option::None);
|
||||
assert!(matches!(absent, std::result::Result::Ok(std::option::Option::None)));
|
||||
let absent = super::canonical_transaction(&reference, fields(&transaction, &meta, &version, &transaction_index, std::option::Option::None));
|
||||
assert!(absent.is_ok());
|
||||
if let std::result::Result::Ok(absent) = absent {
|
||||
assert!(absent.block_time().is_none());
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -242,6 +273,12 @@ fn pre_006_observation_key_is_deterministic_and_endpoint_specific() {
|
||||
let other_endpoint = super::observation_key(&request, &reference, "provider", "endpoint-b");
|
||||
let other_provider = super::observation_key(&request, &reference, "provider-2", "endpoint-a");
|
||||
assert_eq!(first, same);
|
||||
assert_eq!(
|
||||
first.as_bytes(),
|
||||
&[
|
||||
184, 85, 15, 15, 33, 101, 243, 112, 222, 145, 139, 212, 251, 14, 199, 130, 57, 26, 253, 221, 184, 140, 87, 240, 16, 116, 62, 0, 11, 19, 103, 2
|
||||
]
|
||||
);
|
||||
assert_ne!(first, other_endpoint);
|
||||
assert_ne!(first, other_provider);
|
||||
return;
|
||||
@@ -294,6 +331,15 @@ fn pre_006_missing_outcome_contains_only_network_scoped_reference() {
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pre_003_backfill_uses_common_raw_v1_contract_without_changing_frozen_identity() {
|
||||
assert_eq!(crate::RAW_TRANSACTION_FORMAT_ID, ksp_raw_transaction_lib::RAW_TRANSACTION_FORMAT_ID);
|
||||
assert_eq!(crate::RAW_TRANSACTION_FORMAT_VERSION, ksp_raw_transaction_lib::RAW_TRANSACTION_FORMAT_VERSION);
|
||||
assert_eq!(crate::MIN_BACKFILL_SIGNATURE_TEXT_BYTES, ksp_raw_transaction_lib::MIN_RAW_TRANSACTION_SIGNATURE_TEXT_BYTES);
|
||||
assert_eq!(crate::MAX_BACKFILL_SIGNATURE_TEXT_BYTES, ksp_raw_transaction_lib::MAX_RAW_TRANSACTION_SIGNATURE_TEXT_BYTES);
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pre_006_fix_001_raw_acquisition_uses_one_private_indirection() {
|
||||
assert_eq!(std::mem::size_of::<crate::BackfillRawAcquisition>(), std::mem::size_of::<usize>(),);
|
||||
|
||||
Reference in New Issue
Block a user