0.3.16-pre.007

This commit is contained in:
2026-09-22 06:22:40 +02:00
parent a97837b832
commit d768737629
22 changed files with 842 additions and 276 deletions

View File

@@ -1,5 +1,5 @@
<!-- file: crates/ksp-worker-raw-transaction-ingest-lib/README.md -->
<!-- version: 17 -->
<!-- version: 18 -->
# ksp-worker-raw-transaction-ingest-lib
@@ -187,7 +187,7 @@ source hydratante active => quota pending >= 1 et quota in-flight >= 1
Pour garantir simultanément ces bornes et l'absence de starvation structurelle, le démarrage échoue avant spawn si le nombre de sources hydratantes dépasse `admission_queue_capacity` ou `persistence_concurrency`. Le registre des hydrations transactionnelles protège en plus l'ouverture effective des `getTransaction`; Yellowstone Block reste borné par sa partition déterministe et son `JoinSet` possédé.
Les signaux partageant le même `(network, signature, commitment)` sont coalescés cross-source avant le fan-out HTTP. La publication du résultat partagé notifie les followers sous le verrou de registry avant de retirer la clé : une nouvelle génération de leader ne peut donc pas s'intercaler entre retrait et notification. Après canonicalisation, une cache run-local bornée sérialise les acquisitions de même `(network, signature)` : la première passe par l'écriture atomique entity + observation, les suivantes de contenu canonique identique ajoutent uniquement leur observation déterministe. Le Store conserve le guard durable final. En `0.3.15`, il accepte aussi un entrant strictement moins complet lorsque l'unique divergence prouvée est `meta.logMessages` avec exactement un marqueur `Log truncated` après un préfixe identique alors que le canonique stocké n'est pas tronqué ; le canonique complet reste inchangé et l'observation est conservée. Toute autre divergence, ainsi que le sens tronqué -> complet, reste un `content_conflict` terminal. Aucune majorité, préférence provider ou overwrite n'est appliqué.
Les signaux partageant le même `(network, signature, commitment)` sont coalescés cross-source avant le fan-out HTTP. La publication du résultat partagé notifie les followers sous le verrou de registry avant de retirer la clé : une nouvelle génération de leader ne peut donc pas s'intercaler entre retrait et notification. Après canonicalisation, une cache run-local bornée sérialise uniquement les acquisitions de même `(network, signature)` ; elle ne décide jamais du contenu. Chaque acquisition repasse par l'écriture atomique du Store, qui reste l'autorité durable finale. Le comparateur partagé accepte un entrant strictement moins complet lorsque l'unique divergence prouvée est `meta.logMessages` tronqué, promeut atomiquement un entrant prouvé plus complet, et conserve toute divergence conflictuelle/incomparable comme variante durable avec un conflict case ouvert. Aucune majorité, préférence provider ou overwrite n'est appliqué.
Les retries/reroutages HTTP appartiennent à `ksp-onchain-transport-lib`. Le Worker ne possède pas une seconde boucle de retry autour de `getTransaction`.
@@ -283,11 +283,13 @@ entity skipped purged
observation inserted
observation already present
observation not recorded for purged entity
content conflict
variant inserted canonical / exact / compatible less complete
variant promoted compatible more complete
variant quarantined conflict
store failure
```
Un content conflict réel reste terminal en `0.3.15` et n'est jamais converti en succès idempotent. L'exception `logMessages` tronqués décrite ci-dessus est classée avant le conflit comme observation compatible moins complète ; elle ne remplace pas le canonique et n'incrémente pas le terminal `content_conflict`. La conservation de variantes conflictuelles et leur résolution sans arrêt du Worker appartiennent à `0.3.16`.
En `0.3.16-pre.007`, un conflit de contenu classifié par le Store n'est plus une faute terminale : la variante entrante est conservée, le conflict case durable est ouvert ou mis à jour idempotemment et l'outcome `QuarantinedConflict` maintient le Worker `Running` avec une health `Degraded`. Un ancien backend qui renvoie encore `ERROR_CODE_RAW_CONFLICT` reste traité comme erreur terminale de compatibilité. La résolution explicite du conflict case n'appartient pas à cette tranche.
## Snapshots et erreurs

View File

@@ -1,5 +1,5 @@
<!-- file: crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md -->
<!-- version: 15 -->
<!-- version: 16 -->
# Utilisation de ksp-worker-raw-transaction-ingest-lib
@@ -266,7 +266,7 @@ somme de leurs tâches d'hydration actives <= persistence_concurrency
Ces contraintes garantissent au moins une part à chaque source reference-bearing sans introduire de scheduler pondéré. Si les settings ne permettent pas cette répartition, le démarrage échoue avant spawn avec `runtime_invalid`; le caller doit augmenter la capacité concernée ou réduire le nombre de sources nécessitant une hydration.
Pour une même identité `(network, signature)`, le Worker ne choisit ni majorité ni provider préféré et n'écrase pas un contenu divergent ; le Store reste l'autorité durable finale. En `0.3.15`, un canonique complet peut accepter comme observation supplémentaire un entrant identique hors `meta.logMessages` lorsque celui-ci contient exactement un marqueur `Log truncated` après un préfixe identique. Le canonique n'est pas réécrit. Le sens tronqué -> complet et toute autre divergence non prouvée restent des `content_conflict` terminaux ; leur stockage en variantes et leur résolution sont réservés à `0.3.16`.
Pour une même identité `(network, signature)`, le Worker ne choisit ni majorité ni provider préféré et n'écrase pas un contenu divergent ; le Store reste l'autorité durable finale. La cache run-local ne fait que sérialiser cette identité et chaque acquisition atteint le Store. Un entrant prouvé moins complet reste une observation compatible, un entrant prouvé plus complet peut devenir canonique atomiquement, et une divergence conflictuelle/incomparable est conservée comme variante avec un conflict case durable. L'outcome `QuarantinedConflict` n'arrête pas le Worker : il reste `Running` et sa health devient `Degraded`.
## Observer le snapshot concret

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/persistence.rs
// version: 4
// version: 5
/// Canonical entity disposition produced by one Worker Store persistence attempt.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -28,6 +28,7 @@ pub(crate) enum RawTransactionIngestObservationPersistence {
pub(crate) struct RawTransactionIngestPersistenceOutcome {
entity: crate::RawTransactionIngestEntityPersistence,
observation: crate::RawTransactionIngestObservationPersistence,
quarantined_conflict: bool,
}
impl crate::RawTransactionIngestPersistenceOutcome {
@@ -42,6 +43,12 @@ impl crate::RawTransactionIngestPersistenceOutcome {
pub(crate) const fn observation(self) -> crate::RawTransactionIngestObservationPersistence {
return self.observation;
}
/// Reports whether Store durably quarantined divergent or incomparable content without failing the acquisition.
#[must_use]
pub(crate) const fn quarantined_conflict(self) -> bool {
return self.quarantined_conflict;
}
}
/// Private backend-neutral Store persistence port used by the Worker and deterministic tests.
@@ -56,12 +63,6 @@ pub(crate) trait RawTransactionIngestPersistencePort: std::marker::Send + std::m
observation: ksp_store_lib::RawTransactionObservation,
mode: ksp_store_lib::RawTransactionAcquisitionMode,
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawAcquisitionWriteOutcome>>;
/// Records one additional observation for an already durable canonical RAW transaction.
fn record_observation<'a>(
&'a self,
observation: ksp_store_lib::RawTransactionObservation,
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawObservationWriteOutcome>>;
}
impl crate::RawTransactionIngestPersistencePort for ksp_store_lib::Store {
@@ -77,13 +78,6 @@ impl crate::RawTransactionIngestPersistencePort for ksp_store_lib::Store {
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawAcquisitionWriteOutcome>> {
return ksp_store_lib::RawTransactionWrite::persist_raw_transaction_acquisition(self, transaction, observation, mode);
}
fn record_observation<'a>(
&'a self,
observation: ksp_store_lib::RawTransactionObservation,
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawObservationWriteOutcome>> {
return ksp_store_lib::RawTransactionObservationWrite::record_raw_transaction_observation(self, observation);
}
}
#[derive(Clone, Eq, Ord, PartialEq, PartialOrd)]
@@ -92,24 +86,10 @@ struct RawTransactionIngestPersistenceKey {
signature: ksp_store_lib::RawTransactionSignature,
}
#[derive(Clone, Eq, PartialEq)]
struct RawTransactionIngestCanonicalState {
block_time: std::option::Option<ksp_store_lib::RawTimestamp>,
content_hash: ksp_store_lib::RawContentHash,
format_id: ksp_store_lib::RawFormatId,
format_version: u32,
slot: u64,
}
/// Private bounded run-local cache serializing repeated canonical identities before Store writes.
pub(crate) struct RawTransactionIngestPersistenceConvergence {
max_entries: usize,
entries: std::sync::Mutex<
std::collections::BTreeMap<
RawTransactionIngestPersistenceKey,
std::sync::Arc<tokio::sync::Mutex<std::option::Option<RawTransactionIngestCanonicalState>>>,
>,
>,
entries: std::sync::Mutex<std::collections::BTreeMap<RawTransactionIngestPersistenceKey, std::sync::Arc<tokio::sync::Mutex<()>>>>,
}
impl crate::RawTransactionIngestPersistenceConvergence {
@@ -119,10 +99,7 @@ impl crate::RawTransactionIngestPersistenceConvergence {
return Self { max_entries, entries: std::sync::Mutex::new(std::collections::BTreeMap::new()) };
}
fn entry(
&self,
key: RawTransactionIngestPersistenceKey,
) -> std::option::Option<std::sync::Arc<tokio::sync::Mutex<std::option::Option<RawTransactionIngestCanonicalState>>>> {
fn entry(&self, key: RawTransactionIngestPersistenceKey) -> std::option::Option<std::sync::Arc<tokio::sync::Mutex<()>>> {
let mut entries = match self.entries.lock() {
std::result::Result::Ok(value) => value,
std::result::Result::Err(poisoned) => poisoned.into_inner(),
@@ -144,7 +121,7 @@ impl crate::RawTransactionIngestPersistenceConvergence {
if entries.len() >= self.max_entries {
return std::option::Option::None;
}
let entry = std::sync::Arc::new(tokio::sync::Mutex::new(std::option::Option::None));
let entry = std::sync::Arc::new(tokio::sync::Mutex::new(()));
entries.insert(key, std::sync::Arc::clone(&entry));
return std::option::Option::Some(entry);
}
@@ -189,44 +166,14 @@ where
if acquisition.observation().transaction() != &reference {
return std::result::Result::Err(crate::runtime_error("persistence.acquisition_reference_mismatch"));
}
let canonical_state = canonical_state(acquisition.transaction());
let key = RawTransactionIngestPersistenceKey { network: reference.network().clone(), signature: reference.signature() };
let entry = convergence.entry(key);
let entry = match entry {
std::option::Option::Some(value) => value,
std::option::Option::None => return crate::persist_raw_transaction_ingest_acquisition(port, acquisition).await,
};
let mut known_state = entry.lock().await;
if let std::option::Option::Some(known_state_value) = known_state.as_ref() {
if known_state_value != &canonical_state {
return std::result::Result::Err(crate::content_conflict_error());
}
let (_transaction, observation) = acquisition.into_parts();
let result = port.record_observation(observation).await;
return match result {
std::result::Result::Ok(outcome) => map_additional_observation_outcome(outcome),
std::result::Result::Err(error) => map_store_error(error),
};
}
let outcome = crate::persist_raw_transaction_ingest_acquisition(port, acquisition).await;
let outcome = match outcome {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if outcome.entity() != crate::RawTransactionIngestEntityPersistence::SkippedPurged {
*known_state = std::option::Option::Some(canonical_state);
}
return std::result::Result::Ok(outcome);
}
fn canonical_state(transaction: &ksp_store_lib::RawTransaction) -> RawTransactionIngestCanonicalState {
return RawTransactionIngestCanonicalState {
block_time: transaction.block_time(),
content_hash: transaction.payload().content_hash(),
format_id: transaction.payload().format_id().clone(),
format_version: transaction.payload().format_version(),
slot: transaction.slot(),
};
let _identity_guard = entry.lock().await;
return crate::persist_raw_transaction_ingest_acquisition(port, acquisition).await;
}
fn map_store_error(error: ksp_core_lib::Error) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome> {
@@ -239,49 +186,74 @@ fn map_store_error(error: ksp_core_lib::Error) -> ksp_core_lib::Result<crate::Ra
fn map_store_outcome(outcome: ksp_store_lib::RawAcquisitionWriteOutcome) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome> {
let entity = outcome.entity();
let observation = outcome.observation();
let quarantined_conflict = match validate_transaction_variant_outcome(outcome.transaction_variant(), entity, observation) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if entity == ksp_store_lib::RawEntityWriteOutcome::Inserted && observation == ksp_store_lib::RawObservationWriteOutcome::Inserted {
return std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
entity: crate::RawTransactionIngestEntityPersistence::Inserted,
observation: crate::RawTransactionIngestObservationPersistence::Inserted,
quarantined_conflict,
});
}
if entity == ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent && observation == ksp_store_lib::RawObservationWriteOutcome::Inserted {
return std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
entity: crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
observation: crate::RawTransactionIngestObservationPersistence::Inserted,
quarantined_conflict,
});
}
if entity == ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent && observation == ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent {
return std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
entity: crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
observation: crate::RawTransactionIngestObservationPersistence::AlreadyPresent,
quarantined_conflict,
});
}
if entity == ksp_store_lib::RawEntityWriteOutcome::SkippedPurged && observation == ksp_store_lib::RawObservationWriteOutcome::NotRecorded {
return std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
entity: crate::RawTransactionIngestEntityPersistence::SkippedPurged,
observation: crate::RawTransactionIngestObservationPersistence::NotRecorded,
quarantined_conflict,
});
}
return std::result::Result::Err(crate::runtime_error("persistence.store_outcome_invalid"));
}
fn map_additional_observation_outcome(
outcome: ksp_store_lib::RawObservationWriteOutcome,
) -> ksp_core_lib::Result<crate::RawTransactionIngestPersistenceOutcome> {
return match outcome {
ksp_store_lib::RawObservationWriteOutcome::Inserted => std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
entity: crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
observation: crate::RawTransactionIngestObservationPersistence::Inserted,
}),
ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent => std::result::Result::Ok(crate::RawTransactionIngestPersistenceOutcome {
entity: crate::RawTransactionIngestEntityPersistence::AlreadyPresent,
observation: crate::RawTransactionIngestObservationPersistence::AlreadyPresent,
}),
ksp_store_lib::RawObservationWriteOutcome::NotRecorded => {
std::result::Result::Err(crate::runtime_error("persistence.additional_observation_not_recorded"))
fn validate_transaction_variant_outcome(
transaction_variant: std::option::Option<ksp_store_lib::RawTransactionVariantWriteOutcome>,
entity: ksp_store_lib::RawEntityWriteOutcome,
observation: ksp_store_lib::RawObservationWriteOutcome,
) -> ksp_core_lib::Result<bool> {
let recorded_observation =
matches!(observation, ksp_store_lib::RawObservationWriteOutcome::Inserted | ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent);
return match transaction_variant {
std::option::Option::None => std::result::Result::Ok(false),
std::option::Option::Some(ksp_store_lib::RawTransactionVariantWriteOutcome::InsertedCanonical)
if entity == ksp_store_lib::RawEntityWriteOutcome::Inserted && observation == ksp_store_lib::RawObservationWriteOutcome::Inserted =>
{
std::result::Result::Ok(false)
},
_ => std::result::Result::Err(crate::runtime_error("persistence.additional_observation_outcome_invalid")),
std::option::Option::Some(
ksp_store_lib::RawTransactionVariantWriteOutcome::ObservedExact
| ksp_store_lib::RawTransactionVariantWriteOutcome::ObservedCompatibleLessComplete
| ksp_store_lib::RawTransactionVariantWriteOutcome::PromotedCompatibleMoreComplete,
) if entity == ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent && recorded_observation => std::result::Result::Ok(false),
std::option::Option::Some(ksp_store_lib::RawTransactionVariantWriteOutcome::QuarantinedConflict)
if entity == ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent && recorded_observation =>
{
std::result::Result::Ok(true)
},
std::option::Option::Some(ksp_store_lib::RawTransactionVariantWriteOutcome::SkippedPurged)
if entity == ksp_store_lib::RawEntityWriteOutcome::SkippedPurged && observation == ksp_store_lib::RawObservationWriteOutcome::NotRecorded =>
{
std::result::Result::Ok(false)
},
std::option::Option::Some(ksp_store_lib::RawTransactionVariantWriteOutcome::Rehydrated) => {
std::result::Result::Err(crate::runtime_error("persistence.store_variant_rehydrated_in_normal_mode"))
},
_ => std::result::Result::Err(crate::runtime_error("persistence.store_variant_outcome_invalid")),
};
}

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/src/snapshot.rs
// version: 10
// version: 11
/// Runtime-neutral boxed future resolving to one newer concrete RAW transaction ingest Worker snapshot.
pub type RawTransactionIngestSnapshotFuture<'a> =
@@ -1132,6 +1132,13 @@ impl crate::RawTransactionIngestSnapshotPublisher {
let mut entity_skipped_purged_total = self.snapshot.entity_skipped_purged_total;
let mut observation_inserted_total = self.snapshot.observation_inserted_total;
let mut observation_already_present_total = self.snapshot.observation_already_present_total;
let mut content_conflict_total = self.snapshot.content_conflict_total;
if outcome.quarantined_conflict() {
content_conflict_total = match checked_counter(content_conflict_total, "content_conflict_total") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
}
match outcome.entity() {
crate::RawTransactionIngestEntityPersistence::Inserted => {
entity_inserted_total = match checked_counter(entity_inserted_total, "entity_inserted_total") {
@@ -1173,6 +1180,7 @@ impl crate::RawTransactionIngestSnapshotPublisher {
self.snapshot.entity_skipped_purged_total = entity_skipped_purged_total;
self.snapshot.observation_inserted_total = observation_inserted_total;
self.snapshot.observation_already_present_total = observation_already_present_total;
self.snapshot.content_conflict_total = content_conflict_total;
return self.publish(state, admission_queue_depth, in_flight_persistence);
}
@@ -1279,6 +1287,9 @@ fn health_for_state(
return ksp_worker_api::WorkerHealth::Unhealthy;
}
if source_total > 0 && snapshot.source_active == snapshot.source_total {
if snapshot.content_conflict_total > 0 {
return ksp_worker_api::WorkerHealth::Degraded;
}
return ksp_worker_api::WorkerHealth::Healthy;
}
if source_total > 0
@@ -1297,6 +1308,7 @@ fn health_for_state(
ksp_worker_api::WorkerState::Running if snapshot.source_total > 0 && snapshot.source_active < snapshot.source_total => {
ksp_worker_api::WorkerHealth::Degraded
},
ksp_worker_api::WorkerState::Running if snapshot.content_conflict_total > 0 => ksp_worker_api::WorkerHealth::Degraded,
ksp_worker_api::WorkerState::Running => ksp_worker_api::WorkerHealth::Healthy,
ksp_worker_api::WorkerState::Stopping | ksp_worker_api::WorkerState::Stopped => previous,
ksp_worker_api::WorkerState::Faulted(_) => ksp_worker_api::WorkerHealth::Unhealthy,

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/tests/cross_layer_completeness.rs
// version: 3
// version: 4
//! Cross-layer completeness and security canaries through `0.3.14-pre.013`.
@@ -215,7 +215,12 @@ fn v0_3_13_pre_012_all_five_live_source_families_share_one_worker_common_raw_sto
] {
assert!(worker_resources.contains(required), "pre.012 Worker source/pipeline contract missing: {required}");
}
for required in ["persist_raw_transaction_ingest_converged_acquisition", "RawTransactionIngestPersistenceConvergence", "record_observation"] {
for required in [
"persist_raw_transaction_ingest_converged_acquisition",
"RawTransactionIngestPersistenceConvergence",
"let _identity_guard = entry.lock().await",
"persist_raw_transaction_ingest_acquisition(port, acquisition).await",
] {
assert!(worker_persistence.contains(required), "pre.012 Worker persistence contract missing: {required}");
}
for required in ["SolanaStandardWsSession", "HeliusLaserStreamWsSession"] {

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs
// version: 31
// version: 32
//! Dependency firewall canaries for the RAW transaction ingest Worker foundation.
@@ -585,13 +585,16 @@ fn v0_3_13_pre_008_cross_source_convergence_reuses_store_and_transport_facades_o
"RawTransactionIngestPersistenceConvergence",
"persist_raw_transaction_ingest_converged_acquisition",
"RawTransactionIngestPersistencePort",
"record_raw_transaction_observation",
"RawTransactionObservationWrite",
"payload().content_hash()",
"let _identity_guard = entry.lock().await",
"persist_raw_transaction_ingest_acquisition(port, acquisition).await",
"RawTransactionVariantWriteOutcome::QuarantinedConflict",
"content_conflict_error()",
] {
assert!(persistence.contains(required) || root.contains(required), "required pre.008 persistence convergence contract missing: {required}");
}
for forbidden in ["record_raw_transaction_observation", "RawTransactionObservationWrite", "known_state_value", "canonical_state("] {
assert!(!persistence.contains(forbidden), "pre.007 Worker persistence retained stronger run-local arbitration: {forbidden}");
}
for required in [
"RawTransactionIngestGlobalHydrationRegistry",
"fetch_hydration_shared",

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs
// version: 40
// version: 41
//! External public, security, redaction and release-boundary hardening canaries through `v0.3.15-pre.004`.
@@ -764,11 +764,15 @@ fn v0_3_13_pre_008_cross_source_convergence_is_bounded_conflict_checked_and_priv
"max_entries: usize",
"entries.len() >= self.max_entries",
"std::sync::Arc::strong_count(entry) == 1",
"known_state_value != &canonical_state",
"record_observation(observation)",
"persistence.additional_observation_not_recorded",
"let _identity_guard = entry.lock().await",
"persist_raw_transaction_ingest_acquisition(port, acquisition).await",
"outcome.transaction_variant()",
"RawTransactionVariantWriteOutcome::QuarantinedConflict",
] {
assert!(persistence.contains(required), "required pre.008 persistence hardening guard missing: {required}");
assert!(persistence.contains(required), "required persistence serialization/Store-authority guard missing: {required}");
}
for forbidden in ["known_state_value", "canonical_state(", "record_observation(", "RawTransactionObservationWrite"] {
assert!(!persistence.contains(forbidden), "run-local persistence cache still arbitrates before Store: {forbidden}");
}
for required in [
"max_pending: usize",
@@ -803,9 +807,9 @@ fn v0_3_13_pre_009_duplicate_storm_disagreement_and_starvation_guards_are_explic
"tokio::sync::Semaphore::new(max_in_flight)",
"RawTransactionIngestHydrationLeaderGuard",
"source.global_hydration_pending_saturated",
"block_time: transaction.block_time()",
"slot: transaction.slot()",
"content_hash: transaction.payload().content_hash()",
"network: ksp_store_lib::RawNetworkId",
"signature: ksp_store_lib::RawTransactionSignature",
"let _identity_guard = entry.lock().await",
] {
assert!(
admission.contains(required) || persistence.contains(required) || resources.contains(required),

View File

@@ -1,9 +1,10 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/persistence.rs
// version: 3
// version: 4
#[derive(Clone, Copy)]
enum PortResponse {
Outcome(ksp_store_lib::RawEntityWriteOutcome, ksp_store_lib::RawObservationWriteOutcome),
VariantOutcome(ksp_store_lib::RawEntityWriteOutcome, ksp_store_lib::RawObservationWriteOutcome, ksp_store_lib::RawTransactionVariantWriteOutcome),
Conflict,
StoreFailure,
}
@@ -12,8 +13,6 @@ struct FakePersistencePort {
calls: std::sync::atomic::AtomicUsize,
network: ksp_store_lib::RawNetworkId,
normal_mode_seen: std::sync::atomic::AtomicBool,
observation_calls: std::sync::atomic::AtomicUsize,
observation_response: std::option::Option<ksp_store_lib::RawObservationWriteOutcome>,
response: PortResponse,
}
@@ -23,16 +22,9 @@ impl FakePersistencePort {
calls: std::sync::atomic::AtomicUsize::new(0),
network,
normal_mode_seen: std::sync::atomic::AtomicBool::new(false),
observation_calls: std::sync::atomic::AtomicUsize::new(0),
observation_response: std::option::Option::None,
response,
};
}
fn with_observation_response(mut self, observation_response: ksp_store_lib::RawObservationWriteOutcome) -> Self {
self.observation_response = std::option::Option::Some(observation_response);
return self;
}
}
impl crate::RawTransactionIngestPersistencePort for FakePersistencePort {
@@ -52,28 +44,9 @@ impl crate::RawTransactionIngestPersistencePort for FakePersistencePort {
return std::boxed::Box::pin(async move {
return match response {
PortResponse::Outcome(entity, observation) => std::result::Result::Ok(ksp_store_lib::RawAcquisitionWriteOutcome::new(entity, observation)),
PortResponse::Conflict => std::result::Result::Err(ksp_core_lib::Error::new(ksp_store_lib::ERROR_CODE_RAW_CONFLICT, "synthetic conflict")),
PortResponse::StoreFailure => std::result::Result::Err(ksp_core_lib::Error::new(
ksp_core_lib::ErrorCode::new("synthetic_store", "write_failed"),
"synthetic remote failure text that must not escape",
)),
};
});
}
fn record_observation<'a>(
&'a self,
_observation: ksp_store_lib::RawTransactionObservation,
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawObservationWriteOutcome>> {
self.observation_calls.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
let observation_response = self.observation_response;
let response = self.response;
return std::boxed::Box::pin(async move {
if let std::option::Option::Some(outcome) = observation_response {
return std::result::Result::Ok(outcome);
}
return match response {
PortResponse::Outcome(_, observation) => std::result::Result::Ok(observation),
PortResponse::VariantOutcome(entity, observation, transaction_variant) => {
std::result::Result::Ok(ksp_store_lib::RawAcquisitionWriteOutcome::with_transaction_variant(entity, observation, transaction_variant))
},
PortResponse::Conflict => std::result::Result::Err(ksp_core_lib::Error::new(ksp_store_lib::ERROR_CODE_RAW_CONFLICT, "synthetic conflict")),
PortResponse::StoreFailure => std::result::Result::Err(ksp_core_lib::Error::new(
ksp_core_lib::ErrorCode::new("synthetic_store", "write_failed"),
@@ -203,7 +176,36 @@ async fn pre_007_normal_mode_classifies_new_idempotent_and_purged_outcomes() {
}
#[tokio::test(flavor = "current_thread")]
async fn pre_007_content_conflict_maps_to_terminal_worker_code_without_remote_text() {
async fn v0_3_16_pre_007_quarantined_conflict_is_successful_and_explicit() {
let network = match network("mainnet") {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let port = FakePersistencePort::new(
network,
PortResponse::VariantOutcome(
ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent,
ksp_store_lib::RawObservationWriteOutcome::Inserted,
ksp_store_lib::RawTransactionVariantWriteOutcome::QuarantinedConflict,
),
);
let acquisition = match acquisition("mainnet", 8) {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let result = crate::persist_raw_transaction_ingest_acquisition(&port, acquisition).await;
let outcome = match result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return,
};
assert_eq!(outcome.entity(), crate::RawTransactionIngestEntityPersistence::AlreadyPresent);
assert_eq!(outcome.observation(), crate::RawTransactionIngestObservationPersistence::Inserted);
assert!(outcome.quarantined_conflict());
return;
}
#[tokio::test(flavor = "current_thread")]
async fn pre_007_legacy_store_conflict_error_maps_to_terminal_worker_code_without_remote_text() {
let network = match network("mainnet") {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
@@ -303,14 +305,14 @@ async fn pre_007_store_network_mismatch_is_rejected_before_write() {
}
#[tokio::test(flavor = "current_thread")]
async fn v0_3_13_pre_008_same_canonical_identity_uses_atomic_entity_once_then_additional_observation() {
async fn v0_3_16_pre_007_same_identity_is_serialized_but_every_acquisition_reaches_store() {
let network = match network("mainnet") {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let port = FakePersistencePort::new(
network,
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::Inserted),
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent, ksp_store_lib::RawObservationWriteOutcome::Inserted),
);
let convergence = crate::RawTransactionIngestPersistenceConvergence::new(8);
let first = match acquisition_with_observation_key("mainnet", 41, 1) {
@@ -321,56 +323,14 @@ async fn v0_3_13_pre_008_same_canonical_identity_uses_atomic_entity_once_then_ad
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let first_outcome = crate::persist_raw_transaction_ingest_converged_acquisition(&port, first, &convergence).await;
let second_outcome = crate::persist_raw_transaction_ingest_converged_acquisition(&port, second, &convergence).await;
assert!(first_outcome.is_ok());
assert!(second_outcome.is_ok());
assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 1);
assert_eq!(port.observation_calls.load(std::sync::atomic::Ordering::Acquire), 1);
let second_outcome = match second_outcome {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return,
};
assert_eq!(second_outcome.entity(), crate::RawTransactionIngestEntityPersistence::AlreadyPresent);
assert_eq!(second_outcome.observation(), crate::RawTransactionIngestObservationPersistence::Inserted);
assert!(crate::persist_raw_transaction_ingest_converged_acquisition(&port, first, &convergence).await.is_ok());
assert!(crate::persist_raw_transaction_ingest_converged_acquisition(&port, second, &convergence).await.is_ok());
assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 2);
return;
}
#[tokio::test(flavor = "current_thread")]
async fn v0_3_13_pre_008_replayed_same_observation_is_idempotent_without_raw_resubmission() {
let network = match network("mainnet") {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let port = FakePersistencePort::new(
network,
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::Inserted),
)
.with_observation_response(ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent);
let convergence = crate::RawTransactionIngestPersistenceConvergence::new(8);
let first = match acquisition_with_observation_key("mainnet", 42, 3) {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let second = match acquisition_with_observation_key("mainnet", 42, 3) {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let first_result = crate::persist_raw_transaction_ingest_converged_acquisition(&port, first, &convergence).await;
assert!(first_result.is_ok());
let second_result = crate::persist_raw_transaction_ingest_converged_acquisition(&port, second, &convergence).await;
let second_outcome = match second_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return,
};
assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 1);
assert_eq!(port.observation_calls.load(std::sync::atomic::Ordering::Acquire), 1);
assert_eq!(second_outcome.observation(), crate::RawTransactionIngestObservationPersistence::AlreadyPresent);
return;
}
#[tokio::test(flavor = "current_thread")]
async fn v0_3_13_pre_009_source_skew_and_provider_disagreement_fail_before_additional_observation() {
async fn v0_3_16_pre_007_run_local_cache_never_rejects_divergent_shapes_before_store() {
let cases = [
(43_u64, std::option::Option::Some(1_700_000_000_i64), "AQID"),
(42_u64, std::option::Option::Some(1_700_000_001_i64), "AQID"),
@@ -383,7 +343,11 @@ async fn v0_3_13_pre_009_source_skew_and_provider_disagreement_fail_before_addit
};
let port = FakePersistencePort::new(
network,
PortResponse::Outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::Inserted),
PortResponse::VariantOutcome(
ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent,
ksp_store_lib::RawObservationWriteOutcome::Inserted,
ksp_store_lib::RawTransactionVariantWriteOutcome::QuarantinedConflict,
),
);
let convergence = crate::RawTransactionIngestPersistenceConvergence::new(8);
let first = match acquisition_with_shape("mainnet", 51, 10, 42, std::option::Option::Some(1_700_000_000), "AQID") {
@@ -394,15 +358,11 @@ async fn v0_3_13_pre_009_source_skew_and_provider_disagreement_fail_before_addit
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
assert!(crate::persist_raw_transaction_ingest_converged_acquisition(&port, first, &convergence).await.is_ok());
let conflict = crate::persist_raw_transaction_ingest_converged_acquisition(&port, divergent, &convergence).await;
let error = match conflict {
std::result::Result::Ok(_) => return,
std::result::Result::Err(value) => value,
};
assert_eq!(error.code(), crate::ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT);
assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 1);
assert_eq!(port.observation_calls.load(std::sync::atomic::Ordering::Acquire), 0);
let first_outcome = crate::persist_raw_transaction_ingest_converged_acquisition(&port, first, &convergence).await;
let second_outcome = crate::persist_raw_transaction_ingest_converged_acquisition(&port, divergent, &convergence).await;
assert!(matches!(first_outcome, std::result::Result::Ok(value) if value.quarantined_conflict()));
assert!(matches!(second_outcome, std::result::Result::Ok(value) if value.quarantined_conflict()));
assert_eq!(port.calls.load(std::sync::atomic::Ordering::Acquire), 2);
}
return;
}

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs
// version: 12
// version: 13
struct ActiveTaskGuard {
active: std::sync::Arc<std::sync::atomic::AtomicUsize>,
@@ -249,6 +249,7 @@ async fn pre_005_supervisor_reaps_completed_source_task_and_still_joins_live_chi
#[derive(Clone, Copy)]
enum RuntimePortResponse {
Conflict,
QuarantinedConflict,
StoreFailure,
StoreFailureBlocked,
SuccessBlocked,
@@ -337,35 +338,13 @@ impl crate::RawTransactionIngestPersistencePort for RuntimePersistencePort {
));
}
}
if matches!(response, RuntimePortResponse::Conflict) {
if matches!(response, RuntimePortResponse::QuarantinedConflict) {
self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
return std::result::Result::Err(ksp_core_lib::Error::new(ksp_store_lib::ERROR_CODE_RAW_CONFLICT, "synthetic conflict"));
}
self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
return std::result::Result::Err(ksp_core_lib::Error::new(
ksp_core_lib::ErrorCode::new("synthetic_store", "write_failed"),
"synthetic store failure",
));
});
}
fn record_observation<'a>(
&'a self,
_observation: ksp_store_lib::RawTransactionObservation,
) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result<ksp_store_lib::RawObservationWriteOutcome>> {
let response = self.response;
return std::boxed::Box::pin(async move {
if matches!(response, RuntimePortResponse::StoreFailureBlocked | RuntimePortResponse::SuccessBlocked) {
let current = self.active.fetch_add(1, std::sync::atomic::Ordering::AcqRel) + 1;
let _active_guard = RuntimePortActiveGuard { active: &self.active };
self.update_max_active(current);
while !self.released.load(std::sync::atomic::Ordering::Acquire) {
self.notify.notified().await;
}
if matches!(response, RuntimePortResponse::SuccessBlocked) {
self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
return std::result::Result::Ok(ksp_store_lib::RawObservationWriteOutcome::Inserted);
}
return std::result::Result::Ok(ksp_store_lib::RawAcquisitionWriteOutcome::with_transaction_variant(
ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent,
ksp_store_lib::RawObservationWriteOutcome::Inserted,
ksp_store_lib::RawTransactionVariantWriteOutcome::QuarantinedConflict,
));
}
if matches!(response, RuntimePortResponse::Conflict) {
self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
@@ -505,7 +484,50 @@ async fn pre_007_runtime_bounds_in_flight_store_persistence_to_configured_concur
}
#[tokio::test(flavor = "current_thread")]
async fn pre_007_content_conflict_becomes_terminal_after_private_drain() {
async fn v0_3_16_pre_007_quarantined_conflict_keeps_worker_running_and_degraded() {
let settings = match settings_with_persistence_concurrency(1) {
std::option::Option::Some(value) => value,
std::option::Option::None => return,
};
let network = settings.network().clone();
let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::QuarantinedConflict, true));
let runtime_port: super::PersistencePort = port.clone();
let handle = match super::start_foundation_with_port_and_source_spawner(
settings,
tokio::runtime::Handle::current(),
std::option::Option::Some(runtime_port),
move |children, _stop_receiver, admission_sender| {
let _abort_handle = children.spawn(async move {
let ingress = match runtime_ingress(&network, 10) {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")),
};
if admission_sender.send(ingress).await.is_err() {
return std::result::Result::Err(crate::runtime_error("test.source_failed"));
}
return std::result::Result::Ok(());
});
},
) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return,
};
let source = handle.snapshot_source();
assert!(wait_for_completed(&port, 1).await);
let snapshot = source.current();
assert_eq!(snapshot.worker_snapshot().state(), ksp_worker_api::WorkerState::Running);
assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Degraded);
assert_eq!(snapshot.persisted_total(), 1);
assert_eq!(snapshot.content_conflict_total(), 1);
assert_eq!(snapshot.store_failure_total(), 0);
assert!(handle.request_stop());
let terminal = handle.wait_terminal().await;
assert_eq!(terminal, std::result::Result::Ok(ksp_worker_api::WorkerState::Stopped));
return;
}
#[tokio::test(flavor = "current_thread")]
async fn pre_007_legacy_store_conflict_error_becomes_terminal_after_private_drain() {
let settings = match settings_with_persistence_concurrency(1) {
std::option::Option::Some(value) => value,
std::option::Option::None => return,