7.6 KiB
Utilisation de ksp-worker-raw-transaction-ingest-lib
Cette page décrit l'utilisation de la façade publique source-neutral de ksp-worker-raw-transaction-ingest-lib. Le caller possède la composition du runtime Tokio et du Store ; la crate Worker ne construit ni Config, ni backend physique, ni Transport.
Construire les settings
Les identités sont validées par leurs couches propriétaires :
fn worker_settings() -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings> {
let network = match ksp_store_lib::RawNetworkId::new("mainnet") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let worker_id = match ksp_worker_api::WorkerId::new("raw-ingest-mainnet-0001") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
return std::result::Result::Ok(ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings::with_defaults(
network,
worker_id,
));
}
Le Worker kind est fixe :
assert_eq!(
ksp_worker_raw_transaction_ingest_lib::RAW_TRANSACTION_INGEST_WORKER_KIND_CODE,
"raw_transaction_ingest",
);
Pour des limites explicites :
let settings = ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings::new(
network,
worker_id,
512,
16,
std::time::Duration::from_secs(10),
);
Bornes publiques :
admission_queue_capacity 1 ..= 65_536 défaut 256
persistence_concurrency 1 ..= 64 défaut 8
shutdown_drain_timeout 100 ms ..= 30 s défaut 5 s
Une valeur hors borne retourne worker_raw_transaction_ingest.settings_invalid avec uniquement le nom stable du champ invalide.
Démarrer sur le runtime du caller
RawTransactionIngestWorker::start doit être appelé depuis un contexte possédant déjà un runtime Tokio courant. Le Store est passé dans un Arc et doit cibler exactement le même RawNetworkId que les settings.
async fn start_worker(
settings: ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings,
store: std::sync::Arc<ksp_store_lib::Store>,
) -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestHandle> {
return ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestWorker::start(settings, store);
}
Le démarrage retourne avant la fin du Worker. Aucun JoinHandle Tokio n'est exposé au caller.
Un appel hors runtime Tokio courant retourne une erreur worker_raw_transaction_ingest.runtime_invalid. Un mismatch réseau Store/settings est également rejeté avant le spawn.
Observer le snapshot concret
let source = handle.snapshot_source();
let current = source.current();
let sequence = current.worker_snapshot().sequence();
let state = current.worker_snapshot().state();
let queue_depth = current.admission_queue_depth();
let in_flight = current.in_flight_persistence();
Pour attendre une valeur plus récente :
let observed = source.current().worker_snapshot().sequence();
let newer = source.wait_for_change(observed).await;
assert!(newer.worker_snapshot().sequence().is_after(observed));
Le flux est latest-value : les transitions intermédiaires peuvent être coalescées. Ne pas l'utiliser comme journal d'événements exhaustif.
Utiliser la projection Worker API
La même source implémente ksp_worker_api::WorkerSnapshotSource :
let source = handle.worker_snapshot_source();
let current = ksp_worker_api::WorkerSnapshotSource::current(&source);
let observed = current.sequence();
let newer = ksp_worker_api::WorkerSnapshotSource::wait_for_change(&source, observed).await;
assert!(newer.sequence().is_after(observed));
Cette projection commune ne contient que les dimensions génériques Worker. Les compteurs d'admission/persistence restent disponibles sur RawTransactionIngestSnapshot.
Lire les compteurs concrets
Les compteurs principaux sont monotones :
let snapshot = handle.snapshot_source().current();
let admitted = snapshot.admitted_total();
let canonicalized = snapshot.canonicalized_total();
let persisted = snapshot.persisted_total();
let inserted = snapshot.entity_inserted_total();
let existing = snapshot.entity_already_present_total();
let purged = snapshot.entity_skipped_purged_total();
let observations = snapshot.observation_inserted_total();
let conflicts = snapshot.content_conflict_total();
let store_failures = snapshot.store_failure_total();
let source_failures = snapshot.source_failure_total();
let backpressure = snapshot.backpressure_wait_total();
admission_queue_depth() et in_flight_persistence() sont des gauges latest-value, pas des totaux cumulés.
L'épuisement d'un compteur ou de la séquence est terminal et utilise worker_raw_transaction_ingest.counter_exhausted au lieu de wrapper silencieusement.
Demander un stop et attendre le terminal
let accepted = handle.request_stop();
let terminal = handle.wait_terminal().await;
match terminal {
std::result::Result::Ok(state) => {
assert!(state.is_terminal());
}
std::result::Result::Err(error) => return std::result::Result::Err(error),
}
if accepted {
// La première demande de stop a été remise au runtime vivant.
}
request_stop() est idempotent. Il exprime une intention coopérative ; le terminal n'est publié qu'après traitement du drain et des tâches possédées.
Le shutdown est borné par shutdown_drain_timeout. Si les tâches ne peuvent pas être drainées à temps, le Worker les abort puis les rejoint avant de publier worker_raw_transaction_ingest.drain_timeout.
Interpréter les terminaux faultés
Les codes Worker publics sont :
settings_invalid settings techniques hors contrat
runtime_invalid invariant runtime/lifecycle impossible
store_failed erreur Store non-conflict classifiée
content_conflict transaction canonique incompatible avec le contenu durable
counter_exhausted compteur/séquence monotone arrivé à sa borne
source_failed tâche source privée terminée en erreur
drain_timeout drain du shutdown hors deadline
Un consumer doit traiter le ErrorCode comme contrat stable et ne pas dépendre d'un texte backend/provider arbitraire.
Frontière d'admission
La queue d'admission et le type ingress sont volontairement privés. La façade publique actuelle ne propose ni send, ni enqueue, ni registration d'un adapter réseau.
Les adapters d'acquisition qui alimentent cette queue appartiennent à l'implémentation du Worker et doivent conserver :
backpressure mpsc bornée
stop prioritaire sur send bloqué
canonicalisation via ksp-raw-transaction-lib
persistance via ksp-store-lib
aucun backend physique direct
aucun provider dans l'identité canonique
Le caller ne doit pas contourner cette frontière en écrivant directement dans les structures privées du Worker.
Composition supérieure
Le pattern de composition attendu reste :
Config / application / service owner
-> construit le Store
-> choisit ultérieurement les sources/adapters supportés
-> construit RawTransactionIngestSettings
-> démarre RawTransactionIngestWorker
-> observe RawTransactionIngestSnapshotSource
-> demande stop lorsque nécessaire
Le Worker ne reçoit pas de requête historique métier. Une campagne signature/program_id/plage/limite appartient à ksp-job-backfill-lib, pas à ce runtime continu.