v0.3.12-pre.012

This commit is contained in:
2026-09-09 23:24:20 +02:00
parent 474b894d1c
commit 79e278064d
10 changed files with 493 additions and 193 deletions

View File

@@ -1,9 +1,9 @@
<!-- file: crates/ksp-worker-raw-transaction-ingest-lib/USAGE.md -->
<!-- version: 1 -->
<!-- version: 2 -->
# 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.
Cette page décrit la façade publique de `ksp-worker-raw-transaction-ingest-lib`. Le caller possède la composition du runtime Tokio, du `Store` et des ressources Transport ; le Worker ne lit pas Config, ne construit pas un backend physique et ne lit pas de secret depuis l'environnement.
## Construire les settings
@@ -19,20 +19,14 @@ fn worker_settings() -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_l
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,
));
return std::result::Result::Ok(ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestSettings::with_defaults(network, worker_id));
}
```
Le Worker kind est fixe :
```rust
assert_eq!(
ksp_worker_raw_transaction_ingest_lib::RAW_TRANSACTION_INGEST_WORKER_KIND_CODE,
"raw_transaction_ingest",
);
assert_eq!(ksp_worker_raw_transaction_ingest_lib::RAW_TRANSACTION_INGEST_WORKER_KIND_CODE, "raw_transaction_ingest");
```
Pour des limites explicites :
@@ -57,22 +51,85 @@ 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
## Choisir le mode de démarrage
`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.
### Fondation sans source productive
`RawTransactionIngestWorker::start(settings, store)` démarre la fondation runtime sans source Transport. Ce mode reste utile aux tests/compositions qui veulent uniquement le lifecycle, les snapshots et le contrat de shutdown.
```rust
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);
let handle = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestWorker::start(settings, store) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
```
### Source Yellowstone productive
Pour l'ingestion live, le caller compose d'abord les ressources via les crates propriétaires, puis construit :
```rust
fn runtime_resources(
yellowstone_channel: ksp_onchain_transport_lib::YellowstoneGrpcChannel,
subscribe_request: ksp_onchain_transport_lib::YellowstoneSubscribeRequest,
http_pool: ksp_onchain_transport_lib::HttpTransportPool,
hydration_role: ksp_onchain_transport_lib::HttpRoleName,
) -> ksp_core_lib::Result<ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestRuntimeResources> {
let source = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestYellowstoneSource::new(
yellowstone_channel,
subscribe_request,
http_pool,
hydration_role,
) {
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::RawTransactionIngestRuntimeResources::new(source));
}
```
Le démarrage retourne avant la fin du Worker. Aucun `JoinHandle` Tokio n'est exposé au caller.
Puis :
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.
```rust
let handle = match ksp_worker_raw_transaction_ingest_lib::RawTransactionIngestWorker::start_with_runtime_resources(
settings,
store,
runtime_resources,
) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
```
Les deux entrées exigent un runtime Tokio courant et un `Store` portant exactement le même `RawNetworkId` que les settings. Le démarrage avec ressources exige également que la source Yellowstone/HTTP cible ce même réseau.
## Préparer la source Yellowstone
La `YellowstoneSubscribeRequest` doit :
- être valide selon Transport ;
- contenir au moins une famille transaction-bearing admise par le Worker ;
- utiliser explicitement `Confirmed` ou `Finalized` ;
- rester compatible avec le réseau du `YellowstoneGrpcChannel`.
Le `HttpTransportPool` doit posséder au moins une route compatible avec le rôle d'hydration et `getTransaction` sur le même réseau. La construction de `RawTransactionIngestYellowstoneSource` vérifie ces invariants sans ouvrir la connexion réseau.
Le caller ne passe pas de signature, `program_id`, plage de slots ou limite historique au Worker. Ces paramètres appartiennent à un Job Backfill, pas au service continu.
## Comprendre la pipeline live
Les familles Yellowstone sont traitées ainsi :
```text
Transaction -> signal -> HTTP getTransaction -> Common RAW -> admission
TransactionStatus -> signal -> HTTP getTransaction -> Common RAW -> admission
Block -> un signal par transaction -> HTTP getTransaction -> Common RAW -> admission
BlockMeta -> continuity-only
Slot -> continuity-only
Account/Ping/Pong/Entry -> sans RAW Transaction dans cette verticale
```
Les signaux de même `(network, signature, commitment)` sont coalescés avant l'hydration HTTP. Le Worker ne possède pas une boucle de retry HTTP : reroutage/retry/backoff restent dans `ksp-onchain-transport-lib`.
## Observer le snapshot concret
@@ -84,6 +141,10 @@ 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();
let hydration_pending = current.hydration_pending();
let frontier = current.processing_frontier_slot();
let oldest_pending = current.oldest_pending_slot();
let source_state = current.source_state();
```
Pour attendre une valeur plus récente :
@@ -94,7 +155,7 @@ 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.
Le flux est latest-value : les transitions intermédiaires peuvent être coalescées. Ne pas l'utiliser comme journal exhaustif.
## Utiliser la projection Worker API
@@ -108,11 +169,9 @@ let newer = ksp_worker_api::WorkerSnapshotSource::wait_for_change(&source, obser
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`.
La projection commune contient seulement les dimensions génériques Worker. Les compteurs et frontiers spécifiques restent sur `RawTransactionIngestSnapshot`.
## Lire les compteurs concrets
Les compteurs principaux sont monotones :
## Lire les compteurs
```rust
let snapshot = handle.snapshot_source().current();
@@ -128,11 +187,43 @@ 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();
let reconnects = snapshot.source_reconnect_total();
let replay_attempts = snapshot.source_replay_attempt_total();
let proven_gaps = snapshot.source_continuity_gap_total();
```
`admission_queue_depth()` et `in_flight_persistence()` sont des gauges latest-value, pas des totaux cumulés.
`admission_queue_depth()`, `in_flight_persistence()` et `hydration_pending()` sont des gauges latest-value. Les compteurs cumulés ne wrapent jamais silencieusement ; l'épuisement est terminal avec `worker_raw_transaction_ingest.counter_exhausted`.
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.
## Interpréter la processing frontier
`processing_frontier_slot()` est la plus haute slot de travail source réellement observé qui n'est pas bloquée par un pending plus ancien connu. `oldest_pending_slot()` expose ce plus ancien pending lorsqu'il existe.
Cette frontier est strictement run-local :
```text
elle ne prouve pas que toutes les transactions blockchain d'une slot ont été observées
elle ne prouve pas la persistence durable des ingress déjà envoyés à l'admission
elle n'est pas persistée entre deux runs
elle n'est pas un checkpoint Backfill
```
Un `Missing` HTTP règle le signal du point de vue source-processing sans créer de RAW. Un ingress `Available` n'est réglé qu'après envoi réussi vers l'admission centrale. Une hydration abandonnée au shutdown est retirée des pending sans faire avancer artificiellement la frontier.
## Interpréter reconnect et replay
`source_state()` peut retourner `Active`, `Reconnecting`, `Closing`, `Closed` ou `Failed` après démarrage de la source productive.
Les compteurs ont des sémantiques distinctes :
```text
source_reconnect_total reconnects automatiques réussis observés
source_replay_attempt_total tentatives de reconnect portant une demande de replay
source_continuity_gap_total gaps de rétention prouvés par Transport
```
Une tentative de replay n'est pas une preuve de continuité. Le Worker ne choisit pas `from_slot` et ne traite pas directement `SubscribeReplayInfo` ; ces mécanismes appartiennent à Transport.
Si `source_continuity_gap_total` augmente, le Worker fault avec `worker_raw_transaction_ingest.source_failed`. Il ne déclenche pas automatiquement `ksp-job-backfill-lib`.
## Demander un stop et attendre le terminal
@@ -152,55 +243,37 @@ if accepted {
}
```
`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.
`request_stop()` est idempotent. Le terminal n'est publié qu'après le drain borné et la récupération 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 :
## Interpréter les faults
```text
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
content_conflict contenu canonique incompatible avec l'identité 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
source_failed source/replay/hydration terminée par une erreur classifiée
drain_timeout drain de 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 :
```text
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.
Un consumer doit traiter l'`ErrorCode` comme contrat stable et ne pas dépendre d'un texte backend/provider arbitraire.
## Composition supérieure
Le pattern de composition attendu reste :
Le pattern attendu est :
```text
Config / application / service owner
-> résout endpoints, credentials et rôles
-> construit YellowstoneGrpcChannel + YellowstoneSubscribeRequest
-> construit HttpTransportPool + hydration role
-> construit le Store
-> choisit ultérieurement les sources/adapters supportés
-> construit RawTransactionIngestRuntimeResources
-> construit RawTransactionIngestSettings
-> démarre RawTransactionIngestWorker
-> démarre RawTransactionIngestWorker::start_with_runtime_resources
-> 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.
Le Worker ne reçoit pas de requête historique métier et ne dépend pas de Config. Une campagne `signature/program_id/plage/limite` appartient à `ksp-job-backfill-lib`.