v0.3.14-pre.007

This commit is contained in:
2026-09-12 07:47:16 +02:00
parent ae712eae64
commit 887664f3c0
9 changed files with 619 additions and 17 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/continuity.rs
// version: 6
// version: 7
/// Maximum number of slots admitted by one private continuity HTTP discovery window outside this module.
pub(crate) const MAX_RAW_TRANSACTION_INGEST_CONTINUITY_DISCOVERY_WINDOW_SLOTS: u64 = MAX_RAW_TRANSACTION_INGEST_REPAIR_DISCOVERY_WINDOW_SLOTS;
@@ -174,6 +174,43 @@ impl RawTransactionIngestGapRange {
}
}
/// Private run-local obligation for one transaction reference already observed by a live source.
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct RawTransactionIngestKnownReferenceObligation {
commitment: ksp_onchain_transport_lib::SolanaCommitment,
reference: ksp_store_lib::RawTransactionReference,
slot: u64,
}
impl crate::RawTransactionIngestKnownReferenceObligation {
/// Creates one known-reference obligation without inventing any absence or coverage proof.
pub(crate) fn new(
reference: ksp_store_lib::RawTransactionReference,
slot: u64,
commitment: ksp_onchain_transport_lib::SolanaCommitment,
) -> ksp_core_lib::Result<Self> {
if commitment == ksp_onchain_transport_lib::SolanaCommitment::Processed {
return std::result::Result::Err(crate::runtime_error("continuity.known_reference_processed_unsupported"));
}
return std::result::Result::Ok(Self { commitment, reference, slot });
}
/// Returns the exact commitment under which this reference must be resolved.
pub(crate) const fn commitment(&self) -> ksp_onchain_transport_lib::SolanaCommitment {
return self.commitment;
}
/// Returns the already-known canonical transaction reference.
pub(crate) const fn reference(&self) -> &ksp_store_lib::RawTransactionReference {
return &self.reference;
}
/// Returns the already-observed slot associated with the reference.
pub(crate) const fn slot(&self) -> u64 {
return self.slot;
}
}
/// Private run-local anchor for one WebSocket continuity incident.
///
/// The anchor never derives slots from wall-clock time. Its inclusive start is the latest slot actually observed by that source before Transport reported a

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs
// version: 30
// version: 31
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -119,6 +119,8 @@ pub(crate) use self::continuity::RawTransactionIngestContinuityCapabilityDescrip
pub(crate) use self::continuity::RawTransactionIngestContinuityContracts;
/// Private provider-neutral coverage scope used by run-local continuity proof contracts.
pub(crate) use self::continuity::RawTransactionIngestCoverageScope;
/// Private run-local obligation for one transaction reference already observed by a live source.
pub(crate) use self::continuity::RawTransactionIngestKnownReferenceObligation;
/// Private run-local WebSocket incident anchor built only from observed source slots and safe Transport counters.
pub(crate) use self::continuity::RawTransactionIngestWebSocketIncidentAnchor;
/// Creates one terminal content-conflict error without copying conflicting material into diagnostics.

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs
// version: 32
// version: 33
use sha2::Digest; // rust-rules: trait-import
@@ -4026,6 +4026,20 @@ struct RawTransactionIngestHydrationFetch {
observed: RawTransactionIngestObservedTransaction,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum RawTransactionIngestKnownReferenceMissingDisposition {
AwaitCoverage,
PreferBlockSlot,
}
enum RawTransactionIngestKnownReferenceHydrationResolution {
Available(crate::RawTransactionIngress),
Missing {
disposition: RawTransactionIngestKnownReferenceMissingDisposition,
obligation: crate::RawTransactionIngestKnownReferenceObligation,
},
}
#[derive(Clone)]
enum RawTransactionIngestSharedHydrationResult {
Available(RawTransactionIngestObservedTransaction),
@@ -4219,14 +4233,24 @@ impl RawTransactionIngestHydrationCoordinator {
self.pending_signal_count -= pending.signals.len();
for pending_signal in pending.signals {
let signal_slot = pending_signal.signal.slot;
let ingress = finalize_hydration(hydration, settings, pending_signal.signal, pending_signal.received_at, &fetched.observed);
let ingress = match ingress {
let resolution = resolve_known_reference_hydration(hydration, settings, pending_signal.signal, pending_signal.received_at, &fetched.observed);
let resolution = match resolution {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let ingress = match ingress {
std::option::Option::Some(value) => value,
std::option::Option::None => {
let ingress = match resolution {
RawTransactionIngestKnownReferenceHydrationResolution::Available(value) => value,
RawTransactionIngestKnownReferenceHydrationResolution::Missing { disposition, obligation } => {
if obligation.commitment() != hydration.commitment
|| obligation.reference().network() != &hydration.network
|| obligation.slot() != signal_slot
{
return std::result::Result::Err(crate::runtime_error("source.known_reference_obligation_mismatch"));
}
match disposition {
RawTransactionIngestKnownReferenceMissingDisposition::AwaitCoverage
| RawTransactionIngestKnownReferenceMissingDisposition::PreferBlockSlot => {},
}
if let std::result::Result::Err(error) = processing_frontier.settle_pending(signal_slot) {
return std::result::Result::Err(error);
}
@@ -4389,6 +4413,38 @@ fn shared_hydration_result(
};
}
fn resolve_known_reference_hydration(
hydration: &RawTransactionIngestHydrationContext,
settings: &crate::RawTransactionIngestSettings,
signal: RawTransactionIngestSourceSignal,
received_at: ksp_store_lib::RawTimestamp,
observed: &RawTransactionIngestObservedTransaction,
) -> ksp_core_lib::Result<RawTransactionIngestKnownReferenceHydrationResolution> {
let reference = ksp_store_lib::RawTransactionReference::new(signal.network.clone(), signal.signature);
let obligation = match crate::RawTransactionIngestKnownReferenceObligation::new(reference, signal.slot, hydration.commitment) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let ingress = finalize_hydration(hydration, settings, signal, received_at, observed);
let ingress = match ingress {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if let std::option::Option::Some(value) = ingress {
return std::result::Result::Ok(RawTransactionIngestKnownReferenceHydrationResolution::Available(value));
}
let block_slot_supported = match http_role_supports_rpc_method(&hydration.http_pool, &hydration.hydration_role, "getBlock", hydration.network.as_str()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let disposition = if block_slot_supported {
RawTransactionIngestKnownReferenceMissingDisposition::PreferBlockSlot
} else {
RawTransactionIngestKnownReferenceMissingDisposition::AwaitCoverage
};
return std::result::Result::Ok(RawTransactionIngestKnownReferenceHydrationResolution::Missing { disposition, obligation });
}
fn finalize_hydration(
hydration: &RawTransactionIngestHydrationContext,
settings: &crate::RawTransactionIngestSettings,