5.2 KiB
Utilisation de kb-pipeline
Objectif
La crate expose les campagnes bornées d’acquisition, d’extraction, de replay et les inspections stateful nécessaires aux applications et scénarios.
Valider une requête d’extraction Core
fn validate_pending_extraction() -> kb_core::Result<()> {
let request = kb_pipeline::CoreExtractionRequest {
source: kb_pipeline::CoreExtractionSource::Pending,
limit: 1_000,
max_concurrent_extractions: 4,
force_replay: false,
};
request.validate()
}
Une extraction ciblée peut utiliser CoreExtractionSource::Signatures, SlotRange ou ProgramId.
Exécuter une campagne d’extraction Core
async fn run_core_extraction<S, O>(
store: &S,
observer: &O,
request: kb_pipeline::CoreExtractionRequest,
) -> kb_core::Result<kb_pipeline::CoreExtractionSummary>
where
S: kb_store::CanonicalTransactionStore
+ kb_store::CoreExtractionStore
+ Sync,
O: kb_pipeline::CoreExtractionObserver,
{
let result = kb_pipeline::execute_core_extraction(store, observer, request).await;
match result {
Ok(summary) => Ok(summary),
Err(error) => Err(error),
}
}
Préparer une requête de decode replay
fn validate_decode_request(
selection: kb_store::DecodeSelectionFilter,
) -> kb_core::Result<()> {
let request = kb_pipeline::DecodeReplayRequest {
campaign_id: "manual-replay-001".to_string(),
selection,
decoder_names: Vec::new(),
dispatch_policy: kb_pipeline::DecodeDispatchPolicy::HighestPriority,
max_concurrent_inputs: 4,
force_replay: false,
force_replay_all_matching: false,
materialize_after_decode: true,
};
request.validate()
}
Une liste vide dans decoder_names signifie que tous les décodeurs fournis à l’orchestrateur restent éligibles.
Exécuter un decode replay
async fn run_decode_replay<S, O>(
store: &S,
observer: &O,
request: kb_pipeline::DecodeReplayRequest,
decoders: &[&dyn kb_lib::MdApiInstructionDecoder],
materializers: &[&dyn kb_lib::MdApiEventMaterializer],
) -> kb_core::Result<kb_pipeline::DecodeReplaySummary>
where
S: kb_store::DecodeReplayStore + Sync,
O: kb_pipeline::DecodeReplayObserver,
{
let result = kb_pipeline::execute_decode_replay(
store,
observer,
request,
decoders,
materializers,
)
.await;
match result {
Ok(summary) => Ok(summary),
Err(error) => Err(error),
}
}
Backfill HTTP
La surface principale utilise BackfillRequest, BackfillObserver, execute_http_backfill et BackfillSummary.
async fn run_backfill<S, O>(
store: &S,
observer: &O,
request: kb_pipeline::BackfillRequest,
) -> kb_core::Result<kb_pipeline::BackfillSummary>
where
S: kb_store::CanonicalTransactionStore + Sync,
O: kb_pipeline::BackfillObserver,
{
let result = kb_pipeline::execute_http_backfill(store, observer, request).await;
match result {
Ok(summary) => Ok(summary),
Err(error) => Err(error),
}
}
Inspections stateful
Les familles publiques comprennent :
inspect_solana_core_stateful_readiness;- inspections SPL Token et ATA ;
- inspections et corrélations Token-2022 ;
- préflight cryptographique Token-2022 ;
- lecture et matérialisation stateful du registre ElGamal.
async fn inspect_classic_token<S>(
store: &S,
request: kb_pipeline::SplTokenStatefulReadinessRequest,
) -> kb_core::Result<kb_pipeline::SplTokenStatefulReadinessReport>
where
S: kb_store::CanonicalTransactionStore + Sync,
{
let result = kb_pipeline::inspect_spl_token_stateful_readiness(store, request).await;
match result {
Ok(report) => Ok(report),
Err(error) => Err(error),
}
}
Les rapports stateful ne constituent jamais une autorisation implicite d’exécution.
Observateurs
Les campagnes longues exposent des traits d’observation distincts pour le backfill, l’extraction Core et le replay. L’observateur peut publier la progression et participer à l’annulation coopérative selon le contrat concerné.
Erreurs et invariants
- toutes les campagnes sont bornées ;
- la progression persistée ne doit avancer qu’après traitement cohérent ;
- le replay doit rester déterministe pour une même entrée et une même version de pipeline ;
- une matérialisation ne doit pas inventer un état confirmé ;
- les erreurs de transport, stockage, décodage et préflight restent distinguées.
Tests de référence
- tests de frontière contiguë et reprise du backfill ;
- tests d’extraction Core et d’idempotence ;
- tests de decode replay, dispatch et matérialisation ;
- tests stateful SPL Token, ATA et Token-2022 ;
- tests de preuves, préflight cryptographique et postconditions ;
- tests du registre ElGamal fail-closed.
Limites durables
- la crate orchestre les traitements mais ne fournit pas d’interface opérateur ;
- elle ne conserve pas de secret de wallet ;
- elle ne remplace pas les scénarios Devnet et validations explicites de
kb-pipeline-demo-scenarios.